diff --git a/.env.example b/.env.example index 8fb3206..f7cc6ad 100644 --- a/.env.example +++ b/.env.example @@ -123,6 +123,8 @@ CORTEX_CMVK_OPENAI_MODEL=gpt-4o-mini CORTEX_CONTRADICTION_ENABLED=true CORTEX_CONTRADICTION_KAFKA=true CORTEX_SEMANTIC_ENABLED=false +# Decay worker interval when Compose profile `api` is up (seconds). Default 24h. +CORTEX_DECAY_INTERVAL_SECONDS=86400 # ── OpenTelemetry (optional — enable distributed tracing) ───────────────── # Start collector stack: docker compose --profile observability up -d diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 67fb1ac..fd3672c 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -5,9 +5,9 @@ --- -## Current Version: v0.1 — Design Phase -**Status:** Design only — no code written -**Date:** 2026-05-11 +## Current Version: v1.0 — Shipped MVP +**Status:** V1 complete (Phases 0–7). Cortex V2 (memory control plane) starts after Phase 0 audit — see `docs/CORTEX_V2.md` / `docs/CURRENT_STATE.md`. +**Date:** 2026-09-18 **Research foundation:** MAGMA (arXiv:2601.03236), Zep/Graphiti (arXiv:2501.13956), A-MEM (NeurIPS 2025), Field-Theoretic Memory (arXiv:2602.21220), SSGM Framework (arXiv:2603.11768) --- @@ -792,6 +792,7 @@ Done when: | Version | Date | Changes | |---|---|---| +| v1.0 | 2026-09-18 | V1 closeout — shipped MVP status; Phase 8–10 folded into V2 | | v0.1 | 2026-05-11 | Initial architecture design — Session 0 | --- diff --git a/CLAUDE.md b/CLAUDE.md index 273844f..e46d570 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -93,10 +93,10 @@ LinkedIn: linkedin.com/in/abhinaysai-kamineni ## Current Phase -**Phase:** 8 — Outcome tracking + coverage scoring -**Status:** Phase 7 complete — live demo URL + README polish on `main` (set `CORTEX_API_ORIGIN` on Vercel for API-backed search) -**Next phase:** Phase 8 deliverables (coverage scorer, outcome linker) -**Target:** Coverage score on every query + outcome-linked decisions in the graph +**Version:** **V1 closed** (Phases 0–7). Next work is **Cortex V2** — do not start V2 Phase 1 (data foundations) until V2 Phase 0 audit (`docs/CURRENT_STATE.md`) is merged. +**Status:** V1 MVP shipped on `main` — connectors, pipeline, scoring, RBAC, GDPR, dashboard, MCP, live demo URL +**Next:** V2 Phase 0 audit only (#67), then P0 Must Have (#68–#72, #77) +**Target:** Formal V1 closeout docs + honest capability claims; V2 begins with gap report --- @@ -112,9 +112,9 @@ LinkedIn: linkedin.com/in/abhinaysai-kamineni | 5 | Contradiction detector + decay engine | Week 3 | ✅ Done | | 6 | React dashboard + knowledge graph explorer | Week 3-4 | ✅ Done | | 7 | Live demo URL + README polish | Week 4 | ✅ Done | -| 8 | Outcome tracking + coverage scoring | Post-launch | ⏳ | -| 9 | Behavioral mining + elicitation bot | Post-launch | ⏳ | -| 10 | Federated cross-org memory | v2 | ⏳ | +| **V1** | Organizational memory MVP (phases 0–7) | — | ✅ **Closed** | +| V2 | Memory control plane (see `docs/CORTEX_V2.md`) | — | ⏳ Phase 0 next | +| ~~8–10~~ | Coverage / outcomes / elicitation / federation | — | Folded into V2 | --- diff --git a/DECISIONS.md b/DECISIONS.md index 0298d40..bfc12b1 100644 --- a/DECISIONS.md +++ b/DECISIONS.md @@ -23,6 +23,14 @@ Agent picks up OPEN instructions at session start, executes, marks DONE. ## ACTIVE INSTRUCTIONS +### 2026-09-18 — Close V1 then V2 Phase 0 only +Priority: HIGH +Status: OPEN +Detail: +- V1 = phases 0–7. Do not treat coverage/outcomes/meetings as V1 deliverables. +- After V1 closeout merges: execute V2 Phase 0 (#67) only — `docs/CORTEX_V2.md` + `docs/CURRENT_STATE.md`. +- Do **not** start V2 Phase 1 (Claim/Evidence migrations / #68) until Phase 0 is merged. + ### 2026-06-10 — LLM-backed CMVK verifiers (production) Priority: HIGH Status: DONE — `CORTEX_CMVK_BACKEND=openai|ollama` (2026-06-10) @@ -283,6 +291,15 @@ access_policy: { --- +### D-018 — 2026-09-18 — Close V1 at phases 0–7; defer remainder to V2 +**Status:** Active +**Decision:** Cortex **V1** is complete at build phases 0–7 (connectors through live demo URL + README polish). Coverage scoring, outcome linking, meeting connectors, elicitation, and federation are **not** V1 — they move to the Cortex V2 roadmap (memory control plane). V1 closeout includes cache invalidation on graph writes (#33) and scheduled decay-worker (#36), plus README/ARCHITECTURE honesty. +**Rationale:** README overclaimed Phase 8 features as shipped. Closing V1 on what actually works enables an honest V2 Phase 0 audit without pretending coverage/outcomes exist. +**Alternatives rejected:** Keep calling Phase 8 "next" as V1 incomplete; ship stub `coverage_score` just to match README. +**Owner:** Abhinaysai + +--- + ## PENDING DECISIONS (need resolution before build) | # | Decision needed | Options | Deadline | Status | @@ -291,4 +308,4 @@ access_policy: { | P-002 | Local Slack message storage for testing | Real Slack workspace / Slack test fixture files | Before Phase 1 | Open | | P-003 | Neo4j hosting for production | Neo4j AuraDB free tier / Self-hosted EC2 | Before Phase 7 | Open | | P-004 | Dashboard visualization library | D3.js (full control) / React Flow (faster) | Before Phase 6 | Open | -| P-005 | Meeting transcript connector | Recall.ai / AssemblyAI / Whisper self-hosted | Before Phase 8 | Open | +| P-005 | Meeting transcript connector | Recall.ai / AssemblyAI / Whisper self-hosted | Before V2 | Open (see #64) | diff --git a/README.md b/README.md index f854721..70cf0c3 100644 --- a/README.md +++ b/README.md @@ -90,19 +90,19 @@ When any agent touches the payments service, Cortex enriches its context automat | Capability | Description | |---|---| -| **Decision capture** | Extracts structured decisions from Slack, GitHub, Jira, Linear, meetings | -| **Knowledge graph** | Neo4j graph: Decision → Person → System → Exception → Outcome | +| **Decision capture** | Extracts structured decisions from Slack, GitHub, Jira, Linear | +| **Knowledge graph** | Neo4j graph: Decision → Person → System → Exception (+ Outcome schema reserved) | | **Active injection** | Pushes relevant context to agents before they act — not after they ask | | **MCP server** | Native MCP endpoint — any Claude, Cursor, or MCP agent gets memory in one config line | | **Importance scoring** | Filters noise at ingestion — only signal reaches the graph | | **Trust scoring** | Bayesian confidence per memory node — bad inputs don't corrupt memory | | **Contradiction detection** | Flags when new events conflict with existing memory — no silent overwrites | -| **Memory decay** | Old memory compresses and archives on a principled schedule | -| **Coverage scoring** | Per-domain completeness estimate — agents know when memory is thin | +| **Memory decay** | Batch decay engine; scheduled via `decay-worker` when the `api` Compose profile is up | | **RBAC** | Graph-level access control — contractors don't see salary decisions | -| **Outcome tracking** | Links decisions to real metrics — memory becomes self-correcting | | **GDPR erasure** | Cascade delete with audit trail; query cache invalidated per workspace | +**V1 scope note:** Coverage scoring, outcome linking, and meeting connectors are **not shipped in V1** — they are tracked for Cortex V2 (see `docs/CORTEX_V2.md` after Phase 0). Schema stubs for Outcome/coverage indices exist for forward compatibility. + --- ## Production hardening (launch-ready) @@ -111,7 +111,8 @@ These behaviors matter when Cortex runs behind auth in preview or production: | Concern | Behavior | |---|---| -| **GDPR erasure** | `POST /gdpr/erase` bumps a per-workspace Redis cache epoch — stale PII cannot be served from `/query` for up to 60s | +| **GDPR erasure** | `POST /gdpr/erase` bumps a per-workspace Redis cache epoch — stale PII cannot be served from `/query` | +| **Pipeline writes** | Successful graph writes bump the same cache epoch — new decisions are not hidden behind a 60s stale cache | | **Dashboard proxy** | nginx on `:3000` forwards `/gdpr` (and `/query`, `/inject`, …) to the API — same-origin demos work | | **Demo smoke test** | `scripts/demo.sh` sources `.env` and sends `Authorization` when `CORTEX_DEMO_API_KEY` or `CORTEX_API_KEYS` is set | | **Pipeline retries** | Transient Neo4j errors do not commit Kafka offsets or land in DLQ — messages are redelivered | @@ -144,7 +145,7 @@ After deploy, verify end-to-end: `./scripts/verify_free_deploy.sh --api https:// ``` ┌──────────────────────────────────────────────────────────────┐ │ CAPTURE LAYER │ -│ Slack · GitHub · Jira · Linear · Meetings · CI/CD │ +│ Slack · GitHub · Jira · Linear (Meetings / CI/CD → V2) │ │ Real-time event streams via webhooks + OAuth connectors │ └───────────────────────────────┬──────────────────────────────┘ │ Kafka @@ -161,19 +162,18 @@ After deploy, verify end-to-end: `./scripts/verify_free_deploy.sh --api https:// │ Episodic → TimescaleDB what happened and when │ │ Semantic → Qdrant what things mean │ │ Structural → Neo4j relationships + causal chains │ -│ Procedural → Neo4j how things are done │ │ Hot cache → Redis <50ms retrieval for live agents │ └───────────────────────────────┬──────────────────────────────┘ │ ┌───────────────────────────────▼──────────────────────────────┐ │ INTELLIGENCE LAYER │ │ Contradiction detector · Decay engine · Trust scorer │ -│ Coverage scorer · Outcome linker · RBAC enforcer │ +│ RBAC enforcer (Coverage / Outcome linker → V2) │ └───────────────────────────────┬──────────────────────────────┘ │ ┌───────────────────────────────▼──────────────────────────────┐ │ CONTEXT API │ -│ MCP server · REST API · Python SDK · TypeScript SDK │ +│ MCP server · REST API · Python SDK │ │ cortex.query() · cortex.inject() · cortex.remember() │ └──────────────────────────────────────────────────────────────┘ ``` @@ -331,12 +331,12 @@ cortex/ │ ├── jira/ │ └── linear/ ├── extraction/ # Decision extractor, entity resolver, classifier -├── scoring/ # Importance scorer, trust scorer, coverage scorer +├── scoring/ # Importance scorer, trust scorer, CMVK ├── graph/ # Neo4j schema, migrations, Cypher queries │ └── migrations/ # V001__initial_schema.cypher, etc. ├── pipeline/ # Kafka extraction worker (raw → graph) -├── memory/ # Episodic (Timescale) + semantic (Qdrant) helpers -├── intelligence/ # Contradiction detector, decay engine, outcome linker +├── memory/ # Episodic (Timescale) + semantic (Qdrant) + cache epoch +├── intelligence/ # Contradiction detector, decay engine ├── api/ # FastAPI application ├── mcp/ # MCP server (TypeScript) ├── sdk/ # Python client (query, inject, remember) @@ -388,9 +388,9 @@ cortex/ | Phase 5 | Contradiction detector + decay engine | ✅ Shipped | | Phase 6 | React dashboard (Ask, memory map, guide, agent inject) | ✅ Shipped | | Phase 7 | Live demo URL + README polish | ✅ Done (wire `CORTEX_API_ORIGIN` for API-backed search) | -| Phase 8 | Outcome tracking + coverage scoring | ⏳ Post-launch | -| Phase 9 | Elicitation bot (implicit knowledge) | ⏳ Post-launch | -| Phase 10 | Federated cross-org memory | ⏳ v2 | +| **V1** | Phases 0–7 — organizational memory MVP | ✅ **Closed** | +| V2 | Memory control plane (reliability gate, evidence, temporal, firewall) | ⏳ See `docs/CORTEX_V2.md` / issue #78 | +| ~~Phase 8–10~~ | Coverage / outcomes / elicitation / federation | Folded into **Cortex V2** roadmap | **CI:** GitHub Actions runs `pytest` + seed dry-run on push/PR ([`.github/workflows/ci.yml`](.github/workflows/ci.yml)). diff --git a/SESSIONS.md b/SESSIONS.md index 015596b..717ee68 100644 --- a/SESSIONS.md +++ b/SESSIONS.md @@ -601,3 +601,29 @@ 1. Set `CORTEX_API_ORIGIN` on Vercel after Render API deploy (Plan A) 2. Implement coverage scorer (`coverage_score` on `/query`) 3. Implement outcome linker (schema already in V004) + +--- + +## Session — 2026-09-18 — V1 closeout polish +**Duration:** ~1h +**Phase:** V1 close (phases 0–7) before Cortex V2 Phase 0 + +### Built +- Merged PR #59 (demo video / launch announcement removal) +- **#33:** `memory/cache_epoch.py` + `GraphWriter` bumps Redis cache epoch on successful writes +- **#36:** Compose `decay-worker` under `api` profile (`CORTEX_DECAY_INTERVAL_SECONDS`) +- README / ARCHITECTURE honesty — drop overclaimed meetings, coverage, outcomes as shipped +- `docs/V1_RELEASE.md` — formal V1 close; Phase 8–10 folded into V2 +- CLAUDE.md: V1 closed; next is V2 Phase 0 only (no V2 Phase 1 yet) + +### State at end +- V1 MVP formally closed pending this PR merge +- Next: V2 Phase 0 audit only (#67) — CURRENT_STATE.md + CORTEX_V2.md + +### Decisions made +- D-018 — Close V1 at phases 0–7; defer coverage/outcomes/meetings to V2 + +### Next session starts with +1. Merge V1 closeout PR +2. V2 Phase 0 only: `docs/CORTEX_V2.md` + `docs/CURRENT_STATE.md` (#67) +3. Do **not** start Claim/Evidence migrations until Phase 0 merges diff --git a/api/memory.py b/api/memory.py index f6159f5..6464944 100644 --- a/api/memory.py +++ b/api/memory.py @@ -12,6 +12,10 @@ from graph.gdpr import GdprErasureService from graph.query import GraphQueryService +from memory.cache_epoch import ( + bump_workspace_cache_epoch, + workspace_cache_epoch_key, +) from memory.episodic import purge_raw_events_for_person from memory.semantic import ( delete_decision_vectors, @@ -22,8 +26,6 @@ log = structlog.get_logger(__name__) -_WORKSPACE_CACHE_EPOCH_PREFIX = "cortex:ws:" - class MemoryService: """Coordinates graph reads with Redis caching and optional semantic search.""" @@ -54,20 +56,15 @@ def _build_redis_client() -> Any | None: return None def _workspace_cache_epoch(self, workspace_id: str) -> int: - """Per-workspace cache generation — bumped on GDPR erasure.""" + """Per-workspace cache generation — bumped on GDPR erase and graph writes.""" if self._redis is None: return 0 - key = f"{_WORKSPACE_CACHE_EPOCH_PREFIX}{workspace_id}:cache_epoch" - raw = self._redis.get(key) + raw = self._redis.get(workspace_cache_epoch_key(workspace_id)) return int(raw) if raw else 0 def invalidate_workspace_cache(self, workspace_id: str) -> None: - """Drop cached query results for a workspace (e.g. after GDPR erasure).""" - if self._redis is None: - return - key = f"{_WORKSPACE_CACHE_EPOCH_PREFIX}{workspace_id}:cache_epoch" - self._redis.incr(key) - log.info("memory.cache.invalidated", workspace_id=workspace_id) + """Drop cached query results for a workspace (GDPR erase or after writes).""" + bump_workspace_cache_epoch(workspace_id, redis_client=self._redis) def _cache_key(self, prefix: str, payload: dict[str, Any]) -> str: workspace_id = str(payload.get("workspace_id", "")) diff --git a/docker-compose.yml b/docker-compose.yml index d233596..8ea66a1 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -341,6 +341,37 @@ services: condition: service_healthy profiles: ["api"] + decay-worker: + build: + context: . + dockerfile: api/Dockerfile + container_name: cortex-decay-worker + restart: unless-stopped + networks: [cortex] + env_file: .env + environment: + NEO4J_URI: bolt://neo4j:7687 + PYTHONPATH: /app + # Interval between decay batches (seconds). Default 24h for local/demo. + CORTEX_DECAY_INTERVAL_SECONDS: ${CORTEX_DECAY_INTERVAL_SECONDS:-86400} + volumes: + - ./:/app + depends_on: + neo4j: + condition: service_healthy + # Periodic batch: run decay engine then sleep. Uses api image (Python deps). + entrypoint: ["/bin/sh", "-c"] + command: + - | + set -e + INTERVAL="$${CORTEX_DECAY_INTERVAL_SECONDS:-86400}" + echo "cortex-decay-worker: interval=$${INTERVAL}s" + while true; do + python -m intelligence.decay_engine || echo "decay_engine failed; will retry after sleep" + sleep "$${INTERVAL}" + done + profiles: ["api"] + api: build: context: . diff --git a/docs/V1_RELEASE.md b/docs/V1_RELEASE.md new file mode 100644 index 0000000..ce9b54f --- /dev/null +++ b/docs/V1_RELEASE.md @@ -0,0 +1,44 @@ +# Cortex V1 — Release Closeout + +**Status:** Closed +**Date:** 2026-09-18 +**Scope:** Build phases 0–7 (organizational memory MVP) + +## What V1 is + +Cortex V1 is a working **organizational memory** system: + +- Capture decisions from Slack, GitHub, Jira, Linear → Kafka +- Extract structured `DecisionEvent`s → importance / trust / CMVK write gates +- Persist to Neo4j with graph-level RBAC; optional Qdrant semantic merge +- Query / inject / remember via FastAPI + MCP +- Contradiction detection, decay batch job (+ Compose `decay-worker`), GDPR erase +- React dashboard + seeded demo + live Vercel URL + +## What V1 is not (deferred to V2) + +| Claim often seen in older docs | Reality | +|---|---| +| Coverage scoring on `/query` | Schema/UI stub only — V2 / #62 / #76 | +| Outcome linking to metrics | Schema only (V004) — V2 / #63 / #75 | +| Meeting / CI/CD connectors | Not implemented — V2 / #64 | +| Memory Reliability Gate | V2 P0 / #71 | +| Evidence Graph | V2 P0 / #68–#69 | +| Procedural memory / firewall | V2 P1 / #73–#74 | + +## V1 closeout fixes (this release) + +- Redis query-cache epoch bumped on **pipeline graph writes** (not only GDPR) — #33 +- `decay-worker` scheduled service under Compose `api` profile — #36 +- README / ARCHITECTURE honesty pass (no overclaimed coverage/outcomes/meetings) +- Formal V1 closed; Phase 8–10 folded into Cortex V2 roadmap (#78) + +## Ops remaining (not blockers for V1 close) + +- Set `CORTEX_API_ORIGIN` on Vercel after Render/Plan A API deploy (#61) +- Public webhook URLs for live connectors (optional for demos; inject scripts work) + +## Next + +1. **V2 Phase 0 only:** `docs/CORTEX_V2.md` + `docs/CURRENT_STATE.md` (#67) +2. Do **not** start V2 Phase 1 (Claim/Evidence migrations) until Phase 0 merges. diff --git a/graph/writer.py b/graph/writer.py index 308c71d..9029f01 100644 --- a/graph/writer.py +++ b/graph/writer.py @@ -26,6 +26,7 @@ from neo4j import Driver, GraphDatabase from graph.rbac import serialize_access_policy +from memory.cache_epoch import bump_workspace_cache_epoch from scoring.trust_scorer import is_writable from scoring.write_pipeline import assert_scored_for_write from shared.models import IMPORTANCE_DISCARD, DecisionEvent @@ -227,6 +228,9 @@ def write( valid_at=valid_at, ) + # Invalidate Redis query/inject cache so new decisions are visible immediately. + bump_workspace_cache_epoch(decision.workspace_id) + log.info( "graph.write.success", event_id=decision.event_id, diff --git a/memory/cache_epoch.py b/memory/cache_epoch.py new file mode 100644 index 0000000..673f0ca --- /dev/null +++ b/memory/cache_epoch.py @@ -0,0 +1,70 @@ +"""Workspace query-cache epoch — bump on memory mutations so Redis keys miss. + +Used by MemoryService (GDPR erase) and GraphWriter (pipeline writes) so +query/inject caches do not serve stale decisions for up to TTL seconds. +""" + +from __future__ import annotations + +import os +from typing import Any + +import structlog + +log = structlog.get_logger(__name__) + +WORKSPACE_CACHE_EPOCH_PREFIX = "cortex:ws:" + + +def workspace_cache_epoch_key(workspace_id: str) -> str: + """Redis key for the per-workspace cache generation counter.""" + return f"{WORKSPACE_CACHE_EPOCH_PREFIX}{workspace_id}:cache_epoch" + + +def bump_workspace_cache_epoch( + workspace_id: str, + *, + redis_client: Any | None = None, +) -> int | None: + """Increment workspace cache epoch. Returns new epoch or None if Redis unavailable.""" + if not workspace_id: + return None + client = redis_client + if client is None: + client = _connect_redis() + if client is None: + return None + key = workspace_cache_epoch_key(workspace_id) + try: + new_epoch = int(client.incr(key)) + except Exception as exc: + log.warning( + "memory.cache.invalidate_failed", + workspace_id=workspace_id, + error=str(exc), + ) + return None + log.info("memory.cache.invalidated", workspace_id=workspace_id, epoch=new_epoch) + return new_epoch + + +def _connect_redis() -> Any | None: + """Best-effort Redis client from env (same defaults as MemoryService).""" + try: + import redis + except ImportError: + return None + host = os.environ.get("REDIS_HOST", "localhost") + port = int(os.environ.get("REDIS_PORT", "6379")) + try: + client = redis.Redis( + host=host, + port=port, + decode_responses=True, + socket_connect_timeout=0.2, + socket_timeout=0.2, + ) + client.ping() + return client + except Exception: + return None diff --git a/tests/graph/test_writer.py b/tests/graph/test_writer.py index fad958f..2ae0c3a 100644 --- a/tests/graph/test_writer.py +++ b/tests/graph/test_writer.py @@ -57,10 +57,11 @@ def _make_decision( affects: list[str] | None = None, rationale: list[str] | None = None, replaces: str | None = None, + workspace_id: str = WORKSPACE_ID, ) -> DecisionEvent: return DecisionEvent( source_raw_event_id="raw-event-id-001", - workspace_id=WORKSPACE_ID, + workspace_id=workspace_id, event_type="decision", content="We decided to migrate payments to CockroachDB.", made_by=["priya@company.com"] if made_by is None else made_by, @@ -177,16 +178,30 @@ def _make_writer_with_mock_session(self) -> tuple[GraphWriter, MagicMock]: return writer, mock_tx + def _write(self, writer: GraphWriter, decision: DecisionEvent) -> str: + """Write with cache-epoch bump mocked (unit tests do not need Redis).""" + with patch("graph.writer.bump_workspace_cache_epoch"): + return writer.write(decision) + def test_returns_event_id_on_success(self) -> None: writer, _ = self._make_writer_with_mock_session() decision = _make_decision() - result = writer.write(decision) + with patch("graph.writer.bump_workspace_cache_epoch") as bump: + result = writer.write(decision) assert result == decision.event_id + bump.assert_called_once_with(decision.workspace_id) + + def test_invalidates_query_cache_after_write(self) -> None: + writer, _ = self._make_writer_with_mock_session() + decision = _make_decision(workspace_id="ws-cache") + with patch("graph.writer.bump_workspace_cache_epoch") as bump: + writer.write(decision) + bump.assert_called_once_with("ws-cache") def test_decision_node_upsert_called(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision() - writer.write(decision) + self._write(writer, decision) calls = [str(c) for c in mock_tx.run.call_args_list] assert any("MERGE (d:Decision" in c for c in calls) @@ -194,7 +209,7 @@ def test_decision_node_upsert_called(self) -> None: def test_person_upsert_called_for_each_author(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision(made_by=["alice@", "bob@"]) - writer.write(decision) + self._write(writer, decision) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] person_calls = [c for c in cypher_calls if "MERGE (p:Person" in c] @@ -203,7 +218,7 @@ def test_person_upsert_called_for_each_author(self) -> None: def test_system_upsert_called_for_each_system(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision(affects=["payments-service", "auth-service"]) - writer.write(decision) + self._write(writer, decision) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] system_calls = [c for c in cypher_calls if "MERGE (s:System" in c] @@ -212,7 +227,7 @@ def test_system_upsert_called_for_each_system(self) -> None: def test_rationale_upsert_called_for_each_rationale(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision(rationale=["Reason A", "Reason B", "Reason C"]) - writer.write(decision) + self._write(writer, decision) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] rationale_calls = [c for c in cypher_calls if "MERGE (r:Rationale" in c] @@ -221,7 +236,7 @@ def test_rationale_upsert_called_for_each_rationale(self) -> None: def test_supersedes_called_when_replaces_set(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision(replaces="prev-decision-id-001") - writer.write(decision) + self._write(writer, decision) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] assert any("SUPERSEDES" in c for c in cypher_calls) @@ -229,7 +244,7 @@ def test_supersedes_called_when_replaces_set(self) -> None: def test_supersedes_not_called_when_replaces_none(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision(replaces=None) - writer.write(decision) + self._write(writer, decision) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] assert not any("SUPERSEDES" in c for c in cypher_calls) @@ -237,7 +252,7 @@ def test_supersedes_not_called_when_replaces_none(self) -> None: def test_default_access_policy_applied(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision() - writer.write(decision) + self._write(writer, decision) # Find the Decision upsert call and check access_policy is present decision_call_kwargs = mock_tx.run.call_args_list[0][1] @@ -255,7 +270,8 @@ def test_custom_access_policy_propagated(self) -> None: "classification": "confidential", "gdpr_subject": False, } - writer.write(decision, access_policy=custom_policy) + with patch("graph.writer.bump_workspace_cache_epoch"): + writer.write(decision, access_policy=custom_policy) decision_call_kwargs = mock_tx.run.call_args_list[0][1] assert "admin" in str(decision_call_kwargs["access_policy"]) @@ -263,7 +279,7 @@ def test_custom_access_policy_propagated(self) -> None: def test_no_person_calls_when_made_by_empty(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision(made_by=[]) - writer.write(decision) + self._write(writer, decision) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] person_calls = [c for c in cypher_calls if "MERGE (p:Person" in c] @@ -272,7 +288,7 @@ def test_no_person_calls_when_made_by_empty(self) -> None: def test_no_system_calls_when_affects_empty(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() decision = _make_decision(affects=[]) - writer.write(decision) + self._write(writer, decision) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] system_calls = [c for c in cypher_calls if "MERGE (s:System" in c] @@ -281,19 +297,19 @@ def test_no_system_calls_when_affects_empty(self) -> None: def test_minimum_above_threshold_writes_successfully(self) -> None: writer, _ = self._make_writer_with_mock_session() decision = _make_decision(importance_score=IMPORTANCE_DISCARD) - result = writer.write(decision) + result = self._write(writer, decision) assert result == decision.event_id def test_temporal_edges_include_invalid_at_on_rationale(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() - writer.write(_make_decision(rationale=["Because scale"])) + self._write(writer, _make_decision(rationale=["Because scale"])) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] rationale_link = next(c for c in cypher_calls if "HAS_RATIONALE" in c) assert "rel.invalid_at = null" in rationale_link def test_temporal_edges_include_invalid_at_on_supersedes(self) -> None: writer, mock_tx = self._make_writer_with_mock_session() - writer.write(_make_decision(replaces="prev-decision-id-001")) + self._write(writer, _make_decision(replaces="prev-decision-id-001")) cypher_calls = [c[0][0] for c in mock_tx.run.call_args_list] supersedes_link = next(c for c in cypher_calls if "SUPERSEDES" in c) assert "r.invalid_at = null" in supersedes_link diff --git a/tests/memory/test_cache_epoch.py b/tests/memory/test_cache_epoch.py new file mode 100644 index 0000000..10d72a4 --- /dev/null +++ b/tests/memory/test_cache_epoch.py @@ -0,0 +1,33 @@ +"""Unit tests for workspace query-cache epoch helpers.""" + +from __future__ import annotations + +from unittest.mock import MagicMock, patch + +from memory.cache_epoch import ( + bump_workspace_cache_epoch, + workspace_cache_epoch_key, +) + + +def test_workspace_cache_epoch_key() -> None: + assert workspace_cache_epoch_key("ws-1") == "cortex:ws:ws-1:cache_epoch" + + +def test_bump_with_client_increments() -> None: + client = MagicMock() + client.incr.return_value = 3 + epoch = bump_workspace_cache_epoch("ws-1", redis_client=client) + assert epoch == 3 + client.incr.assert_called_once_with("cortex:ws:ws-1:cache_epoch") + + +def test_bump_empty_workspace_noop() -> None: + assert bump_workspace_cache_epoch("", redis_client=MagicMock()) is None + + +def test_bump_connects_when_client_omitted() -> None: + fake = MagicMock() + fake.incr.return_value = 1 + with patch("memory.cache_epoch._connect_redis", return_value=fake): + assert bump_workspace_cache_epoch("ws-x") == 1