Skip to content

About

No description, website, or topics provided.

Resources

Stars

9 stars

Watchers

0 watching

Forks

Repository files navigation

Synapse Cortex

See TESTING.md for test setup, protected behaviors and remaining gaps.

Cognitive backend for the Synapse AI Chat application. A stateless REST API that processes conversational data into a dynamic knowledge graph, enabling personalized long-term memory and intelligent context retrieval for AI assistants.


The Story Behind the Project

This project was built in public and documented through a series of articles that cover the motivation, architecture decisions, and evolution of the system:

  1. My Wife Sent 297 Messages in 15 Days. Not to Me. To the AI I Built Her. The Synapse Story - The origin story, motivation, and how it looks today
  2. Beyond RAG: Building an AI Companion with Deep Memory Using Knowledge Graphs — The first version: why traditional RAG wasn't enough and how knowledge graphs became the foundation for long-term AI memory.
  3. Scaling AI Memory: How I Tamed a 120K Token Prompt with Deterministic GraphRAG — The scaling challenge: how the system evolved to handle growing knowledge graphs without blowing up context windows.
  4. Full Circle: Giving My AI's Knowledge Graph a Notion Interface Using MCP — Closing the loop: connecting the knowledge graph to Notion via MCP for a human-friendly interface to the AI's memory.

If you want to understand the why behind the architecture and design choices in this repo, start with those articles.

Frontend

The chat interface that connects to this backend lives in a separate repository: synapse-chat-ai. Together, they form the complete Synapse system — the frontend handles the conversational UI while this backend powers the knowledge graph and memory layer.


📋 Table of Contents


Overview

Synapse Cortex is a knowledge graph-powered backend designed to give AI chat applications long-term memory capabilities. Instead of treating each conversation in isolation, Synapse Cortex:

  1. Ingests conversational data from chat sessions
  2. Extracts entities, relationships, and facts using LLMs
  3. Stores them in a temporal knowledge graph (Neo4j)
  4. Retrieves relevant context for future conversations
  5. Visualizes the knowledge graph for user exploration and debugging

The system is built on Graphiti, a temporal knowledge graph framework that handles entity resolution, relationship extraction, and temporal invalidation of outdated information.


Core Features

1. 🧠 Knowledge Graph Ingestion

  • Session Processing: Converts chat sessions into structured knowledge (entities, relationships, facts)
  • Entity Resolution: Automatically merges duplicate entities (e.g., "Juan", "Juan Gómez", "JG" → single entity)
  • Temporal Awareness: Tracks when information becomes valid/invalid (e.g., "used to work at X" vs "currently works at Y")
  • Intelligent Filtering: Uses degree-based filtering to exclude low-confidence entities

2. 💬 OpenAI-Compatible Chat Completions

  • Streaming SSE: Server-Sent Events format matching OpenAI's API
  • Google Gemini Backend: Uses Gemini models (flash/pro) for generation
  • System Prompt Injection: Seamlessly injects user knowledge into prompts
  • Async Architecture: Non-blocking streaming with FastAPI

3. 🔍 Smart Context Retrieval (Hydration)

  • Two-Phase Compilation:
    1. Entity Definitions: "What these concepts mean for this user"
    2. Relational Dynamics: "How these concepts interact over time"
  • Cypher-Optimized: Direct Neo4j queries bypass Graphiti's abstraction for read performance
  • Connectivity-Based Ranking: Prioritizes well-connected entities over noise

4. 🧩 GraphRAG - Per-Turn Retrieval-Augmented Generation

  • Hybrid Search: Combines semantic embeddings + BM25 full-text search (RRF fusion) via Graphiti
  • Deduplication: Only injects edges/nodes not already present in the hydrated base prompt
  • Automatic Gating: Skips retrieval when the graph fits entirely in the prompt (is_partial: false)
  • Zero-LLM Overhead: No agent or tool-calling loop — deterministic pipeline with ~1s latency

5. 💾 Gemini Context Caching

  • Explicit CachedContent: Compilation is stored server-side in Gemini as a reusable cached prefix, cutting ~75% of input token cost on repeated turns
  • Transparent Fallback: Stale or missing caches auto-recover by inlining the full compilation in a single retry, no client retry logic needed
  • TTL Refresh on Hit: Successful cache hits extend TTL fire-and-forget so active users never hit expiration mid-session
  • Stateless: The cacheName is returned to the client and forwarded on each chat request — the backend keeps no user→cache map
  • Size-Gated: Compilations below ~1k tokens skip caching (Gemini minimum), falling back to inlining without overhead

6. 🗺️ Knowledge Graph Visualization

  • React-Force-Graph Format: Nodes and links ready for frontend rendering
  • Real-Time Corrections: Natural language memory edits via Graphiti's episode pipeline
  • Temporal Filtering: Only shows valid (non-expired) relationships

7. 📤 Notion Export

  • Graph-to-Notion Pipeline: Exports a user's knowledge graph into structured Notion databases with a summary page
  • Dynamic Schema Design: Gemini analyzes the graph and designs 3-10 category databases with optimal column schemas
  • MCP Agent Integration: Uses the Notion MCP server via create_react_agent for flexible row creation and page building
  • Async with Polling: Fire-and-forget pattern (202 + status polling) with step-level progress tracking
  • Clean Page on Export: Automatically removes all existing content under the parent page before each export
  • Feedback Loop Columns: Every database includes "Needs Review" (checkbox) and "Correction Notes" (rich_text) for corrections
  • Per-Request Auth: Notion token is passed per-request (not server-side) for multi-tenant use

8. 🔄 Notion Correction Import

  • Feedback Loop: Reads user corrections from exported Notion databases and applies them back to the knowledge graph
  • Smart Row Updates: MCP agent intelligently updates affected columns or deletes rows that are no longer relevant
  • Graphiti Integration: Uses add_episode() with custom_extraction_instructions to ensure corrections are language-consistent and properly contradict outdated edges
  • Partial Failure Handling: Individual correction failures don't block the pipeline; detailed failure reports included in the result

9. 🔐 Security & Rate Limiting

  • API Key Authentication: All endpoints (except /health) require X-API-SECRET header
  • Concurrency Control: Configurable semaphore limits to avoid LLM rate limits (429 errors)
  • CORS Middleware: Supports cross-origin requests for web frontends

Technical Architecture

System Overview

┌────────────────────────────────────────────────────────────────────┐
│                         CLIENT APPLICATION                          │
│                     (Synapse AI Chat Frontend)                      │
└───────────────────────────┬────────────────────────────────────────┘
                            │ REST API (JSON/SSE)
                            ▼
┌────────────────────────────────────────────────────────────────────┐
│                         SYNAPSE CORTEX API                          │
│                          (FastAPI + Uvicorn)                        │
├────────────────────────────────────────────────────────────────────┤
│                                                                     │
│  ┌───────────────┐  ┌───────────────┐  ┌───────────────┐          │
│  │  Ingestion    │  │   Hydration   │  │  Generation   │          │
│  │   Service     │  │    Service    │  │    Service    │          │
│  └───────┬───────┘  └───────┬───────┘  └───────┬───────┘          │
│          │                  │                  │                   │
│          └──────────────────┼──────────────────┘                   │
│                             │                                      │
│  ┌──────────────────────────┴───────────────────────────┐          │
│  │              GRAPHITI CORE LAYER                     │          │
│  │  (Entity Resolution, Relationship Extraction,        │          │
│  │   Temporal Management, Embedding Search)             │          │
│  └──────────────────────────┬───────────────────────────┘          │
│                             │                                      │
└─────────────────────────────┼──────────────────────────────────────┘
                              │
                              ▼
         ┌────────────────────────────────────────┐
         │         NEO4J GRAPH DATABASE           │
         │  (Knowledge Graph Storage + Vectors)   │
         └────────────────────────────────────────┘
                              ▲
                              │
         ┌────────────────────────────────────────┐
         │      GOOGLE GEMINI API (External)      │
         │  - LLM: gemini-3-flash-preview         │
         │  - Embeddings: gemini-embedding-001    │
         │  - Reranker: gemini-3-flash-preview    │
         └────────────────────────────────────────┘

Deployment Architectures

Production (Docker Compose + Caddy)

┌─────────────────────────────────────────────────────────────────┐
│                        DIGITAL OCEAN DROPLET                     │
├─────────────────────────────────────────────────────────────────┤
│                                                                  │
│   Internet (HTTPS)                                               │
│        │                                                         │
│        ▼                                                         │
│   ┌──────────────┐                                               │
│   │    Caddy     │  Automatic SSL (Let's Encrypt)               │
│   │  (Port 443)  │  Reverse Proxy + HTTPS Termination           │
│   └──────┬───────┘                                               │
│          │                                                       │
│          ▼                                                       │
│   ┌──────────────┐                                               │
│   │   FastAPI    │  Python 3.12 + Uvicorn                       │
│   │  (Port 8000) │  API Logic + Graphiti Integration            │
│   └──────┬───────┘                                               │
│          │                                                       │
│          ▼                                                       │
│   ┌──────────────┐                                               │
│   │    Neo4j     │  Graph DB + Vector Index                     │
│   │  (Port 7687) │  Persistent Volume (neo4j_data)              │
│   └──────────────┘                                               │
│                                                                  │
└─────────────────────────────────────────────────────────────────┘

Local Development

┌─────────────────────────────────────────────────────────────────┐
│                       DEVELOPER MACHINE                          │
├─────────────────────────────────────────────────────────────────┤
│                                                                  │
│   localhost:8000 ──▶ uvicorn (--reload)                          │
│                                  │                               │
│                                  ▼                               │
│                          FastAPI App (Native)                    │
│                                  │                               │
│                                  ▼                               │
│                          Docker: Neo4j (7687)                    │
│                          Browser UI: http://localhost:7474       │
│                                                                  │
└─────────────────────────────────────────────────────────────────┘

Backend Components

1. Ingestion Service (app/services/ingestion.py)

Purpose: Process chat sessions into the knowledge graph (async fire-and-forget)

Key Functions:

  • accept_session(): Validates, creates job entry, launches background task, returns 202 immediately
  • _process_background(): Runs graphiti.add_episode() asynchronously; updates job store on completion
  • Validates sessions (minimum message count, character threshold)
  • Formats messages into Graphiti episode format

Data Flow:

POST /ingest → Validation → Create Job → 202 Accepted
                    │
                    └─▶ Background: Graphiti.add_episode() → Update job store

GET /ingest/status/{jobId} → Poll until completed → Hydrate on-demand from Neo4j → Return compilation

Validation Rules:

  • Minimum 1 message
  • Minimum 5 total characters

Job Store (app/services/job_store.py): In-memory dict tracks job status. Compilation is fetched from Neo4j on-demand when status is "completed", then the job is cleaned from memory. Requires WEB_CONCURRENCY=1 (single worker).

2. Hydration Service (app/services/hydration.py)

Purpose: Build user knowledge compilations from the graph

Supports two strategies selected via version parameter: V1 (full dump) and V2 (budget-aware).

Architecture (shared):

  • Direct Cypher Queries: Bypasses Graphiti for read performance
  • Degree-Based Filtering: Only includes entities with >= 2 connections (default)
  • Two-Phase Output:
    1. Conceptual Definitions: Entity summaries ordered by connectivity
    2. Relational Dynamics: Relationships with facts, timestamps, and temporal context

Output Format:

#### 1. CONCEPTUAL DEFINITIONS & IDENTITY ####
- **Entity Name**: Summary of what this entity represents
...

#### 2. RELATIONAL DYNAMICS & CAUSALITY ####
- Entity1 relates to Entity2: "fact describing relationship" [since: 2026-02-17]
...

### STATS ###
Definitions: X | Relations: Y | Est. Tokens: ~Z

V2: Budget-Aware Compilation (app/services/hydration_v2.py)

V1 dumps everything into the system prompt with no size control. V2 introduces a cascading waterfill allocation that maximizes context within a character budget (~120k chars / ~30k tokens), ensuring the compilation never blows up the context window.

How it works:

  1. Quality Gates -- Only entities with degree >= 2 and temporally valid edges are fetched (same as V1).

  2. Fast Path -- If the total formatted text fits within the budget, return everything with is_partial: false.

  3. Waterfill Allocation -- When the budget is exceeded:

Budget: 120,000 chars
├── Block A: Nodes (40% = 48,000 chars)
│   Iterate by degree DESC, include atomically until budget exceeded.
│   Unused chars roll over to Block B.
│
└── Block B: Edges (60% + rollover)
    ├── P1: Hub-to-Hub (both nodes in top 20% by degree) — almost always included
    ├── P2: Hub + Recency (one hub node, sorted by date DESC)
    └── P3: Long Tail (low-degree nodes, sorted by date DESC, fills remaining budget)

Hub classification is done entirely in Python using a dictionary built from the node query results -- no extra DB calls. Edges are classified with O(1) lookups against node_degree_map.

Feature flagging: Pass "version": "v2" in the /hydrate request body. Defaults to "v1".

V2 returns metadata for future GraphRAG deduplication:

{
  "compilationMetadata": {
    "is_partial": true,
    "total_estimated_tokens": 29500,
    "included_node_ids": ["uuid-1", "uuid-2"],
    "included_edge_ids": ["uuid-x", "uuid-y"]
  }
}

This lets the frontend persist which nodes/edges are already in the system prompt, so dynamic retrieval only injects what the LLM doesn't already know:

new_edges = [e for e in retrieved_edges if e.id not in metadata.included_edge_ids]

3. GraphRAG Service (app/services/graph_rag.py)

Purpose: Per-turn retrieval-augmented generation from the knowledge graph

When the hydrated compilation is partial (budget-constrained by V2), GraphRAG runs a hybrid search on every chat turn to retrieve long-tail edges and nodes that didn't make the cut, deduplicates them against what's already in the prompt, and injects the new context before generation.

Pipeline:

User message → Build query (last 3 messages) → Graphiti hybrid search (semantic + BM25)
  → Deduplicate against compilationMetadata.included_*_ids
  → Format context block → Append to system message → Stream Gemini

Gating logic (get_rag_skip_reason):

Condition Action
No user_id Skip — no graph to search
No compilationMetadata Skip — no dedup info available
is_partial == false Skip — full graph already in prompt
Otherwise Run GraphRAG search

Key design decisions:

  • No agent/tool-calling: Direct deterministic pipeline — avoids the extra LLM round-trip latency that agent frameworks add (~2-5s vs ~1s)
  • Deduplication by UUID: Uses included_edge_ids / included_node_ids from V2 metadata to ensure zero redundancy with the base prompt
  • Text sanitization: Decodes HTML entities and collapses stray newlines in facts/summaries before injection
  • Graceful degradation: On search failure, proceeds without RAG context (never blocks generation)

Telemetry (OTel rag.* namespace):

  • rag.enabled, rag.skipped_reason
  • rag.search_duration_ms, rag.total_duration_ms
  • rag.raw_edges_count, rag.deduped_edges_count, rag.injected_edges_count
  • rag.raw_nodes_count, rag.deduped_nodes_count, rag.injected_nodes_count
  • rag.query_chars, rag.context_block_chars

4. Generation Service (app/services/generation.py)

Purpose: OpenAI-compatible streaming chat completions

Features:

  • SSE Format: data: {json}\n\n chunks
  • Gemini Integration: Google GenAI SDK with async streaming
  • System Prompt Handling: Prepends system messages to first user message (Gemini requirement)
  • Error Handling: Graceful error streaming with error chunks

Stream Format:

data: {"id": "chatcmpl-...", "choices": [{"delta": {"role": "assistant"}, ...}]}

data: {"id": "chatcmpl-...", "choices": [{"delta": {"content": "Hello"}, ...}]}

data: {"id": "chatcmpl-...", "choices": [{"delta": {}, "finish_reason": "stop"}]}

data: [DONE]

5. Graph Service (app/services/graph.py)

Purpose: Knowledge graph visualization and memory correction

Key Functions:

  • Get Graph: Retrieves nodes and links in react-force-graph format
  • Correct Memory: Applies natural language corrections via Graphiti's episode pipeline

Correction Strategy: Instead of direct CRUD operations (which break embeddings), corrections are processed as new episodes. Graphiti automatically:

  • Invalidates outdated relationships
  • Creates new relationships
  • Maintains temporal integrity

Example: "I no longer want to apply for the O-1 visa, I decided to stay in Colombia" → Graphiti invalidates old edges and creates new ones

6. Notion Export Service (app/services/notion_export.py)

Purpose: Export a user's knowledge graph into structured Notion databases

Pipeline (6 sequential steps, each with its own OTel span):

Step Name Method Description
0 Clean Page Notion SDK (blocks.children.list + DELETE) Remove all existing content under the parent page
1 Hydrate HydrationService.build_user_knowledge(v1) Build the full graph compilation from Neo4j
2a Analyze Schemas Gemini structured output (SchemaResult) Design 3-10 category databases with column schemas
2b Extract Entries Gemini structured output (ExtractionResult) Extract rows for each category (one LLM call per category)
3 Create Databases Notion SDK (notion.request()) Create one Notion database per category under the parent page
4 Populate MCP agent (create_react_agent) Fill databases with extracted rows (batched, 12 rows per agent call)
5 Summarize MCP agent Create a "Knowledge Graph Overview" page with links to all databases

MCP Lifecycle: Steps 4-5 spawn a Node.js subprocess (npx @notionhq/notion-mcp-server) that communicates with Python via stdin/stdout JSON-RPC. The async with stdio_client(...) context manager handles startup, communication, and cleanup. A semaphore (max_concurrent_exports=3) limits concurrent subprocesses.

Job Store (app/services/notion_export_job_store.py): Tracks current_step, categories_count, entries_count, database_ids, and summary_page_url. Same in-memory pattern as the ingest job store.

7. Notion Correction Service (app/services/notion_correction.py)

Purpose: Read user corrections from Notion and apply them back to the knowledge graph

Pipeline (3 steps, each with its own OTel span):

Step Name Method Description
0 Discover Databases Notion SDK (blocks.children.list) Find all child databases under the parent page
1 Scan Flagged Rows Notion SDK (databases/{id}/query) Query each database for rows with "Needs Review" checked
2a Correct Graph graphiti.add_episode() Apply correction with custom_extraction_instructions for language + contradiction handling
2b Update Notion Row MCP agent (create_react_agent) Intelligently update affected columns or delete the row if no longer relevant

Correction Strategy: Each correction is processed as a new Graphiti episode. The custom_extraction_instructions steer the LLM to extract only corrected facts (not the old state), and to write everything in the user's preferred language. Graphiti's built-in contradiction detection automatically invalidates outdated edges.

MCP Agent Decision: The agent receives the full context (current properties, column schema, updated node summaries, new facts, invalidated facts) and chooses to either update the row or archive it.

Job Store (app/services/notion_correction_job_store.py): Tracks current_step, databases_scanned, corrections_found, corrections_applied, corrections_failed, and failed_corrections (per-row error details).


API Endpoints

Authentication

All endpoints (except /health) require an X-API-SECRET header matching SYNAPSE_API_SECRET environment variable.

Endpoints Reference

🟢 GET /health

Purpose: Health check for load balancers and monitoring
Auth: ❌ None required
Response:

{
  "status": "ok",
  "service": "synapse-cortex"
}

📥 POST /ingest

Purpose: Accept a chat session for async processing (fire-and-forget)
Auth: ✅ Required (X-API-SECRET)
Status: 202 Accepted
Request Body:

{
  "jobId": "convex-queue-id-abc123",
  "userId": "user-123",
  "sessionId": "session-abc",
  "messages": [
    {
      "role": "user",
      "content": "I'm planning to move to Spain next year",
      "timestamp": 1704067200000
    },
    {
      "role": "assistant",
      "content": "That's exciting! What's driving this decision?",
      "timestamp": 1704067205000
    }
  ],
  "metadata": {
    "sessionStartedAt": 1704067200000,
    "sessionEndedAt": 1704067300000,
    "messageCount": 2
  }
}

Response (202):

{
  "jobId": "convex-queue-id-abc123",
  "status": "processing"
}

Skipped (insufficient messages): returns immediately with status: "skipped" and userKnowledgeCompilation.

Duplicate submit (same jobId): returns current status without re-processing.

Client Flow: Poll GET /ingest/status/{jobId} until status is "completed" or "failed".


📊 GET /ingest/status/{job_id}

Purpose: Poll for ingest job status and retrieve result when completed
Auth: ✅ Required (X-API-SECRET)

Response (processing):

{
  "jobId": "convex-queue-id-abc123",
  "status": "processing"
}

Response (completed):

{
  "jobId": "convex-queue-id-abc123",
  "status": "completed",
  "userKnowledgeCompilation": "#### 1. CONCEPTUAL DEFINITIONS & IDENTITY ####\n- **Spain**: Country user plans to move to...",
  "metadata": {
    "model": "gemini-3-flash-preview",
    "processing_time_ms": 15000.5,
    "nodes_extracted": 8,
    "edges_extracted": 12,
    "episode_id": "uuid-..."
  }
}

Response (failed):

{
  "jobId": "convex-queue-id-abc123",
  "status": "failed",
  "error": "Error message",
  "code": "GRAPH_PROCESSING_ERROR"
}

404: Job not found (never submitted, already cleaned up, or backend restarted).

Note: When returning a terminal state ("completed" or "failed"), the job is removed from memory. Compilation is hydrated from Neo4j on-demand for completed jobs.


💧 POST /hydrate

Purpose: Fetch current user knowledge compilation without processing new data
Auth: ✅ Required
Request Body:

{
  "userId": "user-123",
  "version": "v2"
}
Field Type Default Description
userId string - Required. User/group ID in the graph
version "v1" | "v2" "v1" Hydration strategy. V2 uses budget-aware waterfill allocation

Response (V1):

{
  "success": true,
  "userKnowledgeCompilation": "#### 1. CONCEPTUAL DEFINITIONS & IDENTITY ####\n..."
}

Response (V2) -- includes compilationMetadata for GraphRAG deduplication:

{
  "success": true,
  "userKnowledgeCompilation": "#### 1. CONCEPTUAL DEFINITIONS & IDENTITY ####\n...",
  "compilationMetadata": {
    "is_partial": true,
    "total_estimated_tokens": 29500,
    "included_node_ids": ["uuid-1", "uuid-2"],
    "included_edge_ids": ["uuid-x", "uuid-y"]
  }
}

Use Cases:

  • Debugging current graph state
  • Fetching context without re-indexing
  • Feature-flagging V1 vs V2 compilation from the frontend

💬 POST /v1/chat/completions

Purpose: OpenAI-compatible streaming chat completions
Auth: ✅ Required
Request Body:

{
  "messages": [
    {
      "role": "system",
      "content": "You are a helpful assistant."
    },
    {
      "role": "user",
      "content": "What's the weather like today?"
    }
  ],
  "model": "gemini-3-flash-preview",
  "stream": true
}

Response: Server-Sent Events (SSE)

data: {"id":"chatcmpl-abc123","choices":[{"delta":{"role":"assistant"},...}]}

data: {"id":"chatcmpl-abc123","choices":[{"delta":{"content":"Today"},...}]}

data: {"id":"chatcmpl-abc123","choices":[{"finish_reason":"stop",...}],"usage":{"grounding_enabled":true,"grounding_used":true,"grounding_query_count":1,"grounding_source_count":3,"grounding_support_count":2,"grounding_search_entry_point":"<div>...</div>","grounding_sources":[{"title":"example.org","uri":"https://..."}]}}

data: [DONE]

Google Search grounding is available by default. Gemini decides whether to search for each request. Set GROUNDING_ENABLED=false to remove the tool from chat generations. The final usage chunk distinguishes availability from actual use and includes source metadata without exposing the generated search queries. When present, clients must render grounding_search_entry_point unchanged with the grounded response to display Google's required Search Suggestions.


🗺️ GET /v1/graph/{group_id}

Purpose: Retrieve knowledge graph in react-force-graph format
Auth: ✅ Required
Response:

{
  "nodes": [
    {
      "id": "uuid-1234",
      "name": "Spain",
      "val": 5,
      "summary": "Country user plans to relocate to in 2025"
    }
  ],
  "links": [
    {
      "source": "uuid-1234",
      "target": "uuid-5678",
      "label": "RELATES_TO",
      "fact": "User plans to move to Spain for work opportunities"
    }
  ]
}

Features:

  • Node val = connection count (controls visual sizing)
  • Only returns valid (non-expired) relationships
  • Excludes episodic nodes

✏️ POST /v1/graph/correction

Purpose: Apply natural language memory corrections
Auth: ✅ Required
Request Body:

{
  "group_id": "user-123",
  "correction_text": "I've changed my mind. I no longer want to move to Spain. I'm staying in Colombia."
}

Response:

{
  "success": true
}

How It Works:

  • Correction text is processed as a new Graphiti episode
  • Graphiti automatically invalidates outdated edges
  • Creates new edges reflecting the correction
  • Preserves embeddings and temporal integrity

📤 POST /v1/notion/export

Purpose: Export a user's knowledge graph into Notion databases (async) Auth: ✅ Required (X-API-SECRET) Status: 202 Accepted Request Body:

{
  "userId": "user-123",
  "notionToken": "ntn_your_notion_integration_secret",
  "pageName": "Synapse",
  "language": "English"
}
Field Type Default Description
userId string - Required. User/group ID whose graph to export
notionToken string - Required. Notion internal integration secret
pageName string - Required. Name of the parent Notion page (must be shared with the integration)
language string "English" Output language for all generated Notion content

Response (202):

{
  "jobId": "a1b2c3d4-e5f6-7890-abcd-ef1234567890",
  "status": "processing",
  "pageId": "12345678-abcd-1234-abcd-123456789abc"
}

The Notion token and page name are validated synchronously before returning 202. If the token is invalid or the page is not found, you get a 400 error immediately.

Client Flow: Poll GET /v1/notion/export/status/{jobId} for progress and results.


📊 GET /v1/notion/export/status/{job_id}

Purpose: Poll for Notion export job status and retrieve result when completed Auth: ✅ Required (X-API-SECRET)

Response (processing):

{
  "jobId": "a1b2c3d4-...",
  "status": "processing",
  "progress": {
    "currentStep": "populating",
    "categoriesDesigned": 6,
    "entriesExtracted": 42
  }
}

Pipeline steps reported via currentStep: "hydrating" → "analyzing" → "extracting_entries" → "creating_databases" → "populating" → "summarizing" → "done".

Response (completed):

{
  "jobId": "a1b2c3d4-...",
  "status": "completed",
  "result": {
    "databaseIds": {
      "People": "db-id-1",
      "Health & Medications": "db-id-2",
      "Projects & Goals": "db-id-3"
    },
    "summaryPageUrl": "https://notion.so/abc123...",
    "categoriesCount": 6,
    "entriesCount": 42,
    "durationMs": 185000.5
  }
}

Response (failed):

{
  "jobId": "a1b2c3d4-...",
  "status": "failed",
  "error": "Error message",
  "code": "UPSTREAM_TIMEOUT"
}

404: Job not found (never submitted, already consumed, or backend restarted).

Note: Terminal states ("completed" or "failed") remove the job from memory after the response is returned.


🔄 POST /v1/notion/corrections

Purpose: Import user corrections from Notion databases back into the knowledge graph Auth: ✅ Required (X-API-SECRET) Status: 202 Accepted Request Body:

{
  "userId": "user-123",
  "notionToken": "ntn_your_token_here",
  "pageName": "Synapse",
  "language": "Spanish"
}

Response:

{
  "jobId": "d1e2f3a4-...",
  "status": "processing",
  "pageId": "resolved-page-id"
}

Client Flow: Poll GET /v1/notion/corrections/status/{jobId} for progress and results.


📊 GET /v1/notion/corrections/status/{job_id}

Purpose: Poll for Notion correction import job status Auth: ✅ Required (X-API-SECRET)

Response (processing):

{
  "jobId": "d1e2f3a4-...",
  "status": "processing",
  "progress": {
    "currentStep": "applying",
    "databasesScanned": 7,
    "correctionsFound": 3,
    "correctionsApplied": 1,
    "correctionsFailed": 0
  }
}

Pipeline steps reported via currentStep: "scanning" → "applying" → "done".

Response (completed):

{
  "jobId": "d1e2f3a4-...",
  "status": "completed",
  "result": {
    "correctionsFound": 3,
    "correctionsApplied": 2,
    "correctionsFailed": 1,
    "failedCorrections": [
      {"category": "Medications", "title": "Aspirin", "error": "LLM rate limit exceeded"}
    ],
    "durationMs": 95000.5
  }
}

Response (failed):

{
  "jobId": "d1e2f3a4-...",
  "status": "failed",
  "error": "No databases found under the specified Notion page.",
  "code": "NO_DATABASES"
}

Note: Terminal states ("completed" or "failed") remove the job from memory after the response is returned.


Gemini Context Caching

Overview

Every chat turn sends the user's compiled knowledge (~30k tokens) as part of the prompt. Without caching, each request re-bills those tokens at full price. Gemini's explicit context caching stores the compilation server-side as a reusable prefix, so only the short per-turn delta (system date, user message, RAG context) is billed at full price — typically ~75% cheaper on repeated tokens and materially lower TTFT because the cached prefix is pre-tokenized.

This implementation is stateless on the Cortex side: the cacheName is returned to the client after creation, the client persists it, and forwards it on every POST /v1/chat/completions. Cortex does not hold a user→cache map in memory.

Why explicit caching (not implicit)

Google GenAI offers two caching modes:

Mode Trigger Min size Hit guarantee Chosen?
Implicit Automatic on prompt prefix match — Best-effort ❌
Explicit Manual caches.create with resource name ~1k tokens Guaranteed when cached_content is passed ✅

Explicit caching gives deterministic cost savings and hit rates, which is critical when the 30k-token compilation dominates the input budget.

Lifecycle

┌──────────────────────────────────────────────────────────────────┐
│                      CACHE CREATION                               │
│                                                                   │
│  POST /hydrate  or  GET /ingest/status/{jobId} (completed)        │
│       │                                                           │
│       ▼                                                           │
│  HydrationService.build_user_knowledge()                          │
│       │                                                           │
│       ▼                                                           │
│  CacheManager.create_compilation_cache(userId, compilation)       │
│       │                                                           │
│       ├─▶ len(compilation) < 16000 chars? → skip, return None     │
│       │                                                           │
│       └─▶ genai.caches.create(                                    │
│              system_instruction=compilation,                      │
│              tools=[google_search] (when grounding is enabled),   │
│              ttl="900s",                                          │
│              model=<chat_model>,                                 │
│            ) → "cachedContents/xyz..."                           │
│                                                                   │
│  Response includes cacheName — client persists it alongside       │
│  the user's compilation metadata.                                 │
└──────────────────────────────────────────────────────────────────┘

┌──────────────────────────────────────────────────────────────────┐
│                       CACHE USAGE                                 │
│                                                                   │
│  POST /v1/chat/completions { messages, cache_name, compilation }  │
│       │                                                           │
│       ▼                                                           │
│  GenerationService.stream_chat_completion()                       │
│       │                                                           │
│       ├─▶ cache_name set? → call Gemini with                      │
│       │     config=GenerateContentConfig(cached_content=<name>)   │
│       │     (Google Search is already declared in the cache)      │
│       │     (compilation NOT inlined — it lives in the cache)     │
│       │                                                           │
│       ▼                                                           │
│  Peek first chunk under try/except                                │
│       │                                                           │
│       ├─▶ Success: stream chunks, track usage_metadata            │
│       │      │                                                   │
│       │      └─▶ On final chunk: if cached_content_token_count>0 │
│       │           fire-and-forget CacheManager.refresh_ttl()     │
│       │                                                          │
│       └─▶ Cache error (404, expired, model mismatch):            │
│            ├─▶ CacheManager.invalidate_by_name(cache_name)       │
│            ├─▶ Rebuild contents with compilation INLINED         │
│            └─▶ Retry once — always succeeds (same payload as     │
│                pre-cache era)                                     │
└──────────────────────────────────────────────────────────────────┘

Fallback behavior

Caches can become stale in several ways: TTL expired, cache deleted upstream, the cache's model doesn't match the request's model. The generation service treats all of these uniformly:

  1. Peek the first stream chunk under try/except — the Gemini SDK is lazy, so cache errors surface when we pull the first chunk, not when we open the stream.
  2. Detect cache errors with a string-match allowlist (cachedcontent, cache expired, does not match the model in the cached content, etc.).
  3. Invalidate the stale cache name (caches.delete) so no one else hits it.
  4. Rebuild the request contents with inline_compilation=True and retry once. This retry uses the exact payload we would have sent if caching had never been enabled — it always succeeds under the same upstream conditions.
  5. Signal the fallback to the client via usage.cache_fallback_triggered = true in the final SSE chunk. The client uses this to schedule a background re-hydrate so the next turn gets a fresh cache.

The user never sees an error — the fallback adds ~500ms (one extra peek + retry) to the first affected turn and resolves itself by the next one.

TTL refresh on hit

Gemini caches expire by wallclock, not usage. Without active extension, a user in an active conversation would hit expiration at the 15-minute mark and force a fallback.

After every successful cache hit (cached_content_token_count > 0), the generation service spawns a fire-and-forget caches.update(ttl="900s") task. This pushes the expiration window forward, so actively-used caches live as long as the user keeps typing (typical in-session gap: 4-7 min).

The task is held in a module-level _background_tasks set (Python's asyncio only keeps weak references to tasks — without a strong ref, the GC could drop the task mid-execution).

Client responsibilities

The Cortex server is stateless regarding cache ownership. The client (e.g. synapse-chat-ai) must:

  1. Persist cacheName returned by /hydrate or /ingest/status/{jobId}.
  2. Forward it on every chat request as cache_name (alongside compilation — the server needs both: cache for the happy path, compilation for the fallback).
  3. Refresh after fallback — when the SSE response's final usage chunk has cache_fallback_triggered: true, schedule a /hydrate call so the next turn has a valid cache.
  4. Clear on user/persona change — caches are scoped to a user's compilation; switching contexts means the old cacheName is meaningless.

Constraints

Constraint Value Reason
Minimum size 16000 chars (~4000 tokens) Below this, storage cost exceeds savings from cheaper reads
Default TTL 900s (15m) Covers typical in-session gap; minimizes post-session idle storage
Model binding Cache model must equal request model Gemini returns 400 INVALID_ARGUMENT otherwise — caught and fallback-triggered
Auth Full Vertex AI (GCP_PROJECT + credentials) or AI Studio (GOOGLE_API_KEY) Vertex Express (VERTEX_API_KEY) does NOT support caching

The chat model setting (CHAT_MODEL) must match what clients send in request.model. If they diverge, every chat request triggers the model-mismatch fallback — functionally correct but defeats the cache savings.

Observability

The cache.* attribute namespace in Axiom exposes cache decisions and outcomes:

Attribute Description
cache.enabled Whether this request was sent with a cache_name
cache.name The cachedContents/... resource name used
cache.skip_reason Why caching was skipped (no_compilation_in_request, no_cache_name_from_client, compilation_too_small)
cache.hit True when Gemini reported cached_content_token_count > 0
cache.hit_ratio cached_tokens / prompt_tokens
cache.fallback_triggered True when the cache errored and we retried with inlined compilation
cache.fallback_error Truncated error message that triggered the fallback
cache.compilation_chars Size of the compilation that would have been inlined
cache.creation_duration_ms Latency of caches.create (on /hydrate and /ingest/status)

Cache hit ratio over time:

['synapse-cortex-traces']
| where name == 'chat.completion.stream'
| where ['attributes.cache.enabled'] == true
| summarize
    requests = count(),
    avg_hit_ratio = avg(['attributes.cache.hit_ratio']),
    hit_rate = (countif(['attributes.cache.hit'] == true) * 100.0 / count()),
    fallback_rate = (countif(['attributes.cache.fallback_triggered'] == true) * 100.0 / count())
  by bin_auto(_time)

Token savings estimate from caching:

['synapse-cortex-traces']
| where name == 'chat.completion.stream'
| where ['attributes.cache.hit'] == true
| summarize
    cached_tokens = sum(['attributes.chat.tokens.cached']),
    total_prompt_tokens = sum(['attributes.chat.tokens.prompt']),
    savings_pct = (sum(['attributes.chat.tokens.cached']) * 100.0 / sum(['attributes.chat.tokens.prompt']))
  by model = ['attributes.chat.model']

Fallback events (each line = one cache-error recovery):

['synapse-cortex-traces']
| where name == 'chat.completion.stream'
| where ['attributes.cache.fallback_triggered'] == true
| project
    _time,
    cache_name = ['attributes.cache.name'],
    error = ['attributes.cache.fallback_error'],
    prompt_tokens = ['attributes.chat.tokens.prompt']
| order by _time desc

Notion Export

Overview

The Notion Export feature lets you export a user's entire knowledge graph into structured, browsable Notion databases. It reads the graph compilation from Neo4j, uses Gemini to dynamically design database schemas and extract entries, then creates everything in Notion under a parent page you specify.

How It Works

1. CLIENT
   └─▶ POST /v1/notion/export
       {userId, notionToken, pageName, language}

2. ROUTE HANDLER (synchronous validation)
   ├─▶ Validate notionToken by resolving pageName → pageId
   ├─▶ Create job in memory store (status: "processing")
   ├─▶ Launch background task (asyncio.create_task)
   └─▶ Return 202 {jobId, pageId}

3. BACKGROUND PIPELINE (NotionExportService)
   │
   ├─▶ Step 0: CLEAN PAGE
   │   └─▶ List all child blocks under the parent page
   │       └─▶ Delete each block (databases, summaries, etc.)
   │           → Ensures a fresh page for every export
   │
   ├─▶ Step 1: HYDRATE
   │   └─▶ HydrationService.build_user_knowledge(userId, v1)
   │       → Full graph compilation text from Neo4j
   │
   ├─▶ Step 2a: DESIGN SCHEMAS
   │   └─▶ Gemini structured output (SchemaResult)
   │       → 3-10 categories with column definitions
   │
   ├─▶ Step 2b: EXTRACT ENTRIES
   │   └─▶ Gemini structured output (ExtractionResult) × N categories
   │       → All rows for each category
   │
   ├─▶ Step 3: CREATE DATABASES
   │   └─▶ Notion SDK (notion.request) × N categories
   │       → One database per category under the parent page
   │       → Each includes "Needs Review" + "Correction Notes" columns
   │
   ├─▶ Step 4: POPULATE (MCP agent)
   │   └─▶ npx @notionhq/notion-mcp-server (stdio subprocess)
   │       └─▶ ReAct agent creates rows via API-post-page
   │           (batched, 12 rows per agent call)
   │
   └─▶ Step 5: SUMMARIZE (MCP agent)
       └─▶ ReAct agent creates "Knowledge Graph Overview" page
           with overview text, database links, and feedback instructions

4. CLIENT (polling loop)
   └─▶ GET /v1/notion/export/status/{jobId}
       ├─▶ {status: "processing", progress: {currentStep, ...}} → poll again
       ├─▶ {status: "completed", result: {databaseIds, summaryPageUrl, ...}}
       └─▶ {status: "failed", error, code}

Notion Setup

  1. Create an internal integration at https://www.notion.so/profile/integrations
  2. Copy the integration secret (starts with ntn_)
  3. Create a parent page in Notion (e.g. "Synapse")
  4. Share the page with your integration (page menu → "Connect to" → your integration)
  5. Use the page name as the pageName parameter in the API request

Example: Full Export Flow

# 1. Start the export
curl -X POST http://localhost:8000/v1/notion/export \
  -H "Content-Type: application/json" \
  -H "X-API-SECRET: your_secret" \
  -d '{
    "userId": "user-123",
    "notionToken": "ntn_your_token_here",
    "pageName": "Synapse",
    "language": "English"
  }'
# → 202 {"jobId": "abc-123", "status": "processing", "pageId": "..."}

# 2. Poll for status (repeat until completed/failed)
curl http://localhost:8000/v1/notion/export/status/abc-123 \
  -H "X-API-SECRET: your_secret"
# → {"jobId": "abc-123", "status": "processing", "progress": {"currentStep": "populating", ...}}

# 3. Final result
# → {"jobId": "abc-123", "status": "completed", "result": {"databaseIds": {...}, "summaryPageUrl": "https://notion.so/...", ...}}

Prerequisites

The Notion export pipeline spawns a Node.js MCP subprocess. The server environment needs:

  • Node.js (v18+) and npx available on PATH
  • Network access to registry.npmjs.org (first run downloads @notionhq/notion-mcp-server)

Axiom Observability

The pipeline emits spans under the export.* attribute namespace:

['synapse-cortex-traces']
| where startswith(name, 'notion_export.')
| summarize
    total = count(),
    avg_duration = avg(['attributes.export.duration_ms']),
    avg_categories = avg(['attributes.export.categories_count']),
    avg_entries = avg(['attributes.export.entries_count']),
    failed = countif(['attributes.operation.status'] == 'failed')
  by bin_auto(_time)

Per-step latency breakdown:

['synapse-cortex-traces']
| where startswith(name, 'notion_export.')
| summarize
    avg_ms = avg(['attributes.duration_ms']),
    p95_ms = percentile(['attributes.duration_ms'], 95),
    calls = count()
  by name
| order by avg_ms desc

Notion Correction Import

Overview

The Notion Correction Import feature closes the feedback loop: users flag incorrect data in the exported Notion databases, and the system reads those corrections, applies them to the knowledge graph, and intelligently updates or deletes the Notion rows.

Each exported database includes two feedback columns:

  • Needs Review (checkbox): Flag a row that needs correction
  • Correction Notes (rich_text): Describe what needs to be fixed

How It Works

1. CLIENT
   └─▶ POST /v1/notion/corrections
       {userId, notionToken, pageName, language}

2. ROUTE HANDLER (synchronous validation)
   ├─▶ Validate notionToken by resolving pageName → pageId
   ├─▶ Create job in memory store (status: "processing")
   ├─▶ Launch background task (asyncio.create_task)
   └─▶ Return 202 {jobId, pageId}

3. BACKGROUND PIPELINE (NotionCorrectionService)
   │
   ├─▶ Step 0: DISCOVER DATABASES
   │   └─▶ List child_database blocks under the parent page
   │       → Map of category name → database ID
   │
   ├─▶ Step 1: SCAN FOR FLAGGED ROWS
   │   └─▶ Query each database filtering "Needs Review" == true
   │       └─▶ Extract property values, types, and correction notes
   │           → List of CorrectionItems
   │
   └─▶ Step 2: APPLY CORRECTIONS (per row, sequentially)
       │
       ├─▶ 2a. CORRECT GRAPH (Graphiti)
       │   └─▶ graphiti.add_episode() with:
       │       - Episode body containing old properties + user correction
       │       - custom_extraction_instructions for language + correction handling
       │       → Returns AddEpisodeResults (updated nodes, edges, invalidated facts)
       │
       └─▶ 2b. UPDATE NOTION ROW (MCP agent)
           └─▶ LangGraph ReAct agent with Notion MCP tools decides:
               - OPTION A: Update row (patch affected properties, uncheck flag)
               - OPTION B: Delete row (archive if entity is no longer relevant)
               Agent receives: current properties, column schema, node summaries,
               new facts, invalidated facts, and user correction notes.

4. CLIENT (polling loop)
   └─▶ GET /v1/notion/corrections/status/{jobId}
       ├─▶ {status: "processing", progress: {currentStep, correctionsApplied, ...}}
       ├─▶ {status: "completed", result: {correctionsFound, correctionsApplied, correctionsFailed, ...}}
       └─▶ {status: "failed", error, code}

Example: Full Correction Flow

# 1. Start the correction import
curl -X POST http://localhost:8000/v1/notion/corrections \
  -H "Content-Type: application/json" \
  -H "X-API-SECRET: your_secret" \
  -d '{
    "userId": "user-123",
    "notionToken": "ntn_your_token_here",
    "pageName": "Synapse",
    "language": "Spanish"
  }'
# → 202 {"jobId": "def-456", "status": "processing", "pageId": "..."}

# 2. Poll for status
curl http://localhost:8000/v1/notion/corrections/status/def-456 \
  -H "X-API-SECRET: your_secret"
# → {"status": "processing", "progress": {"currentStep": "applying", "correctionsApplied": 1, ...}}

# 3. Final result
# → {"status": "completed", "result": {"correctionsFound": 3, "correctionsApplied": 2, "correctionsFailed": 1, ...}}

How the MCP Agent Decides

The agent receives the full context of each correction and makes an intelligent decision:

  • Update: If the correction modifies specific facts (e.g., "the dosage changed to 20mg"), the agent patches the relevant columns while keeping unaffected properties intact
  • Delete: If the correction invalidates the entity entirely (e.g., "this concept is no longer relevant"), the agent archives the row

The language parameter ensures all updated property values are written in the user's preferred language.

Axiom Observability

The correction pipeline emits spans under the correction.* attribute namespace:

['synapse-cortex-traces']
| where startswith(name, 'notion_correction.')
| summarize
    total = count(),
    avg_duration = avg(['attributes.duration_ms']),
    avg_found = avg(['attributes.correction.found']),
    avg_applied = avg(['attributes.correction.applied']),
    avg_failed = avg(['attributes.correction.failed'])
  by bin_auto(_time)

Per-step breakdown (graph correction vs Notion row update):

['synapse-cortex-traces']
| where startswith(name, 'notion_correction.')
| summarize
    avg_ms = avg(['attributes.duration_ms']),
    p95_ms = percentile(['attributes.duration_ms'], 95),
    calls = count()
  by name
| order by avg_ms desc

Data Flow

Ingestion Pipeline (Async with Polling)

1. CLIENT (Convex)
   └─▶ POST /ingest
       {jobId, userId, sessionId, messages, metadata}

2. INGESTION SERVICE (accept_session)
   ├─▶ Validate session (min messages, min chars)
   ├─▶ If insufficient: return 202 {status: "skipped", userKnowledgeCompilation} (hydrate immediately)
   ├─▶ If duplicate jobId: return 202 {status: "processing"} (no re-process)
   ├─▶ Create job in memory store (status: "processing")
   ├─▶ asyncio.create_task(_process_background)
   └─▶ Return 202 {jobId, status: "processing"}

3. BACKGROUND TASK (_process_background)
   └─▶ GRAPHITI CORE
       ├─▶ Extract entities (e.g., "hiking", "User")
       ├─▶ Extract relationships (e.g., User ENJOYS hiking)
       ├─▶ Generate embeddings (gemini-embedding-001)
       ├─▶ Search for existing similar entities
       ├─▶ Merge duplicates (entity resolution)
       ├─▶ Rerank candidates (Gemini reranker)
       └─▶ Write to Neo4j (nodes + edges + timestamps)
   └─▶ Update job store: status "completed" + metadata (no compilation stored)

4. CLIENT (Polling loop, e.g. exponential backoff: 2m, 5m, 10m, 30m)
   └─▶ GET /ingest/status/{jobId}
       ├─▶ 404: job not found → re-submit POST /ingest
       ├─▶ {status: "processing"} → poll again
       └─▶ {status: "completed"} → HYDRATION SERVICE
           ├─▶ Hydrate on-demand from Neo4j (~1-2s)
           ├─▶ Remove job from memory
           └─▶ Return full result to client

Chat Completion Pipeline

1. CLIENT
   └─▶ POST /v1/chat/completions
       {messages, model, stream: true, user_id, compilationMetadata}

2. GRAPH RAG (conditional)
   ├─▶ Gate check: user_id present? compilationMetadata.is_partial == true?
   ├─▶ Build search query from last 3 messages
   ├─▶ Graphiti hybrid search (semantic + BM25, limit=10)
   ├─▶ Deduplicate against included_edge_ids / included_node_ids
   └─▶ Append context block to system message

3. GENERATION SERVICE
   ├─▶ Convert messages to Gemini format
   │   - Extract system prompt (now includes RAG context)
   │   - Prepend to first user message
   │   - Map roles: user/assistant → user/model
   └─▶ Call gemini.generate_content_stream() with Google Search available

4. GEMINI API
   └─▶ Stream token chunks

5. GENERATION SERVICE
   ├─▶ Wrap chunks in OpenAI SSE format
   │   data: {"choices": [{"delta": {"content": "..."}}]}
   ├─▶ Include RAG stats in final usage chunk (rag_enabled, rag_edges, etc.)
   ├─▶ Include grounding decision, counts, and sources in final usage chunk
   └─▶ Send [DONE] marker

6. CLIENT
   └─▶ Receives streamed response

Memory Correction Pipeline

1. CLIENT
   └─▶ POST /v1/graph/correction
       {group_id, correction_text}

2. GRAPH SERVICE
   └─▶ Call graphiti.add_episode(correction_text)

3. GRAPHITI CORE
   ├─▶ Process correction as new episode
   ├─▶ Extract new facts/relationships
   ├─▶ Search for conflicting edges
   ├─▶ Set invalid_at timestamp on outdated edges
   └─▶ Create new edges with valid_at = now

4. NEO4J
   └─▶ Graph updated with temporal integrity

5. CLIENT
   └─▶ Next hydration/graph fetch reflects changes

Technology Stack

Core Framework

  • FastAPI: Modern async web framework with automatic OpenAPI docs
  • Uvicorn: ASGI server with WebSocket and SSE support
  • Pydantic: Data validation and settings management

Knowledge Graph

  • Neo4j: Graph database with vector search capabilities
  • Graphiti Core: Temporal knowledge graph framework with entity resolution
  • Google Gemini: LLM for entity extraction, embeddings, and reranking
    • LLM Model: gemini-3-flash-preview (configurable)
    • Embedding Model: gemini-embedding-001 (3072 dimensions)
    • Reranker Model: gemini-3-flash-preview

Deployment

  • Docker & Docker Compose: Containerization and orchestration
  • Caddy: Automatic HTTPS/SSL with Let's Encrypt
  • Digital Ocean: Cloud hosting (production environment)

Dependencies

fastapi>=0.109.0
uvicorn[standard]>=0.27.0
neo4j>=5.17.0
graphiti-core[google-genai]>=0.27.1
google-genai>=1.0.0
pydantic-settings>=2.1.0
python-dotenv>=1.0.0
sse-starlette>=1.8.0

Observability in Axiom

The API now emits structured OpenTelemetry attributes for request-level and service-level debugging in Axiom.

Attribute namespaces

  • chat.*: completions, token usage, stream metrics, upstream Gemini details
  • rag.*: GraphRAG gating, search latency, edge/node counts, dedup stats
  • ingest.*: job lifecycle, counts, processing metadata
  • hydrate.*: hydration request context and output size
  • export.*: Notion export pipeline steps, categories, entries, database counts
  • correction.*: Notion correction import pipeline, corrections found/applied/failed
  • graph.*: graph retrieval/correction context and result counts
  • db.*: Neo4j query type, records returned, query latency
  • upstream.*: upstream status/error hints for Gemini and HTTP calls
  • grounding.*: Google Search availability, actual use, query/source/support counts
  • error.*: normalized category/code/type/message for failures

Query ideas (Axiom)

Error rate by endpoint:

['synapse-cortex-traces']
| where ['attributes.http.route'] != ''
| summarize
    total = count(),
    failed = countif(['attributes.operation.status'] == 'failed'),
    error_rate = (countif(['attributes.operation.status'] == 'failed') * 100.0 / count())
  by route = ['attributes.http.route']
| order by error_rate desc

Chat completion performance + token usage:

['synapse-cortex-traces']
| where name == 'chat.completion.stream'
| summarize
    requests = count(),
    p95_ms = percentile(['attributes.chat.total_duration_ms'], 95),
    avg_total_tokens = avg(['attributes.chat.tokens.total']),
    avg_prompt_tokens = avg(['attributes.chat.tokens.prompt']),
    avg_completion_tokens = avg(['attributes.chat.tokens.completion'])
  by model = ['attributes.chat.model']
| order by requests desc

Gemini/upstream failures by category/status:

['synapse-cortex-traces']
| where ['attributes.upstream.error_type'] != ''
| summarize
    failures = count()
  by category = ['attributes.error.category'],
     code = ['attributes.error.code'],
     status = ['attributes.upstream.status_code']
| order by failures desc

GraphRAG retrieval performance:

['synapse-cortex-traces']
| where ['attributes.rag.enabled'] == true
| summarize
    requests = count(),
    p95_search_ms = percentile(['attributes.rag.search_duration_ms'], 95),
    avg_injected_edges = avg(['attributes.rag.injected_edges_count']),
    avg_injected_nodes = avg(['attributes.rag.injected_nodes_count']),
    avg_deduped = avg(['attributes.rag.deduped_edges_count'] + ['attributes.rag.deduped_nodes_count']),
    zero_injection_rate = (countif(['attributes.rag.injected_edges_count'] + ['attributes.rag.injected_nodes_count'] == 0) * 100.0 / count())
| order by requests desc

Slow Neo4j queries:

['synapse-cortex-traces']
| where startswith(name, 'db.cypher.')
| where ['attributes.db.query_duration_ms'] > 500
| project
    timestamp,
    name,
    query_type = ['attributes.db.query_type'],
    duration_ms = ['attributes.db.query_duration_ms'],
    records = ['attributes.db.records_returned']
| order by duration_ms desc

Setup & Deployment

Prerequisites

  • Python: 3.12+
  • Node.js: 18+ with npx on PATH (required for the Notion export MCP subprocess)
  • Docker: Latest stable version
  • Docker Compose: V2+
  • Google Gemini API Key: Get one here

Environment Variables

Create a .env file based on .env.example:

Variable Description Default Required
NEO4J_URI Neo4j Bolt connection URI bolt://localhost:7687 Yes
NEO4J_USER Neo4j username neo4j Yes
NEO4J_PASSWORD Neo4j password - Yes
GOOGLE_API_KEY Google Gemini API key - Yes
GRAPHITI_MODEL Gemini model for Graphiti gemini-3-flash-preview No
GROUNDING_ENABLED Make Google Search available to chat generations true No
SYNAPSE_API_SECRET API authentication secret - Yes
SEMAPHORE_LIMIT Max concurrent LLM operations 3 No

Local Development

1. Start Neo4j

docker-compose -f docker-compose.local.yml up -d

Neo4j Access:

  • Browser UI: http://localhost:7474
  • Bolt Protocol: bolt://localhost:7687
  • Credentials: neo4j / localpassword

2. Install Dependencies

python -m venv venv
source venv/bin/activate  # Windows: venv\Scripts\activate
pip install -r requirements.txt

3. Configure Environment

cp .env.example .env
nano .env  # Edit with your credentials

4. Run FastAPI

uvicorn app.main:app --reload

API available at: http://localhost:8000

5. Test Endpoints

# Health check
curl http://localhost:8000/health

# Hydrate (requires X-API-SECRET header)
curl -X POST http://localhost:8000/hydrate \
  -H "Content-Type: application/json" \
  -H "X-API-SECRET: your_secret" \
  -d '{"userId": "test-user"}'

6. Rails-style console (optional)

Interactive Python shell with Graphiti and all project services pre-loaded (like rails c):

# From project root, with venv active
python scripts/console.py

You get a REPL with graphiti, neo4j_driver, graph_service, hydration_service, ingestion_service, generation_service, and settings in scope. Try searches, call services, or run any Python. Install IPython for top-level await and a nicer prompt: pip install ipython.

Examples in the console:

edges = await graphiti.search("what are my preferences?", group_id="user-123")
g = await graph_service.get_graph("user-123")

For a minimal search-only loop (no free-form code), use python scripts/graphiti_repl.py instead.

7. Stop Services

# Stop Neo4j
docker-compose -f docker-compose.local.yml down

# Remove data volume (optional)
docker-compose -f docker-compose.local.yml down -v

Production Deployment

1. Server Provisioning

  • Create a Digital Ocean droplet (Ubuntu 22.04 LTS recommended)
  • Configure firewall to allow ports 80, 443, 22
  • SSH into the server

2. Clone Repository

git clone <your-repo-url> synapse-cortex
cd synapse-cortex

3. Configure Environment

cp .env.example .env
nano .env

Production .env example:

NEO4J_URI=bolt://neo4j:7687
NEO4J_USER=neo4j
NEO4J_PASSWORD=<strong-random-password>
GOOGLE_API_KEY=<your-gemini-api-key>
SYNAPSE_API_SECRET=<your-api-secret>
GRAPHITI_MODEL=gemini-3-flash-preview
SEMAPHORE_LIMIT=3

4. Configure DNS

Create an A record pointing to your droplet's IP:

  • Host: synapse-cortex
  • Value: Droplet IP address
  • Domain: juandago.dev

Result: synapse-cortex.juandago.dev → Droplet IP

5. Deploy Stack

docker-compose up -d --build

This starts:

  • Caddy: Automatic SSL on port 443
  • FastAPI: API server on port 8000 (internal)
  • Neo4j: Graph database on port 7687 (internal)

6. Verify Deployment

# Check containers
docker-compose ps

# View logs
docker-compose logs -f

# Test API
curl https://synapse-cortex.juandago.dev/health

Expected response:

{"status":"ok","service":"synapse-cortex"}

7. Update Deployment

git pull
docker-compose up -d --build

8. Monitor Logs

# All services
docker-compose logs -f

# Specific service
docker-compose logs -f api
docker-compose logs -f neo4j
docker-compose logs -f caddy

9. Backup Neo4j Data

# Stop services
docker-compose down

# Backup neo4j_data volume
docker run --rm -v synapse-cortex_neo4j_data:/data -v $(pwd):/backup ubuntu tar czf /backup/neo4j-backup-$(date +%Y%m%d).tar.gz /data

# Restart services
docker-compose up -d

10. Restore Neo4j Data

# Stop services
docker-compose down

# Restore from backup
docker run --rm -v synapse-cortex_neo4j_data:/data -v $(pwd):/backup ubuntu tar xzf /backup/neo4j-backup-YYYYMMDD.tar.gz -C /

docker run --rm \
  -v $(pwd):/backups \
  -v synapse-cortex_neo4j_data:/data \
  neo4j:5-community \
  neo4j-admin database load neo4j --from-path=/backups --overwrite-destination=true


# Restart services
docker-compose up -d

Demo User Seeding

Pre-populates a demo account with a realistic knowledge graph so visitors can explore the app with an already-formed graph, existing conversation history, and features like Notion export — without running any live conversations.

The token cost is paid once (ingesting conversations through Graphiti). Subsequent resets are free (direct Neo4j batch insert, no LLM calls).

Files

scripts/
  seed_demo.json              # Source conversations (generated in AI Studio, do not modify)
  ingest_demo.py              # Step 1: runs Graphiti on the conversations → builds graph
  export_demo_graph.py        # Step 2: exports graph snapshot → seed_data/demo_graph.json
  reset_demo.py               # Step 3: resets a group_id from the snapshot (repeatable)
  seed_data/
    demo_graph.json           # Committed graph snapshot with DEMO_PLACEHOLDER group_id

How it works

The group_id (Clerk user ID) is dynamic per environment, so the snapshot stores DEMO_PLACEHOLDER as the group_id value. At reset time, reset_demo.py substitutes it with the actual Clerk user ID before inserting.

Step 1 — Ingest (once, costs tokens)

Reads seed_demo.json and calls graphiti.add_episode() for each session. Takes several minutes depending on session count and rate limits.

python scripts/ingest_demo.py
# Uses group_id: demo_seed_YYYYMMDD by default

python scripts/ingest_demo.py --group-id my_custom_group

Step 2 — Export snapshot (once, free)

Exports all Entity/Episodic nodes and RELATES_TO/MENTIONS edges into scripts/seed_data/demo_graph.json. Use --delete-after to clean up the temporary group_id from Neo4j.

python scripts/export_demo_graph.py --group-id demo_seed_YYYYMMDD --delete-after

Commit scripts/seed_data/demo_graph.json — this is the reusable snapshot.

Step 3 — Reset (free, repeatable)

Deletes existing graph data for the demo account and re-inserts from the snapshot. Run this whenever you want to restore the demo to its initial state.

# Dry run to preview what will happen
python scripts/reset_demo.py --group-id <clerk_user_id> --dry-run

# Apply
python scripts/reset_demo.py --group-id <clerk_user_id>

Verify

python scripts/console.py
# In the REPL:
# g = await graph_service.get_graph("<clerk_user_id>")
# len(g.nodes), len(g.links)

Updating the demo data

To replace the demo content entirely:

  1. Edit or regenerate scripts/seed_demo.json
  2. Re-run Steps 1–2 with a new temporary group_id
  3. Commit the updated scripts/seed_data/demo_graph.json

Alternative chat generation with OpenRouter

Set OPEN_ROUTER_API_KEY to enable the four explicit non-Google chat models configured in app/services/openrouter_generation.py. Docker Compose forwards this variable to the API container. Requests to /v1/chat/completions select the alternative with provider: "openrouter" and an allowed model ID. Omitting provider keeps the existing Gemini generation path.

Both paths use the existing GraphRAG step. OpenRouter receives persona instructions, compiled knowledge and retrieved memories as text, plus the original message history. Gemini cache IDs are not sent to OpenRouter and its Google Search tool is not enabled there. Graphiti, embeddings, ingestion and the Vertex generation service retain their existing clients. OpenRouter errors return SSE error events; there is no model fallback. The active choices are GPT-6.1 Sol, Claude Sonnet 5.5, Qwen3.8 Max 0902 and Kimi K2.6; all accept images. Removed model IDs are rejected. Google models continue to use Vertex. Convex resolves the user assignment and freezes it per turn; model identities and usage stay in internal metadata/logs. When PostHog is enabled, OpenRouter emits $ai_generation events for successful and failed calls with trace/session IDs, token usage, reported USD cost and latency in seconds. Analytics failures do not interrupt chat generation.

Run python -m unittest discover -s tests -v to verify the shared retrieval step, provider dispatch, reasoning configuration, image capabilities and SSE handling.

About

No description, website, or topics provided.

Resources

Stars

9 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages