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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -243,7 +243,7 @@ Set `OTEL_EXPORTER_OTLP_ENDPOINT` to enable distributed tracing (OpenTelemetry O
| `GET` | `/decisions/by-system/{system_id}` | Recent decisions affecting a service |
| `GET` | `/decisions/{id}/chain` | SUPERSEDES / trigger lineage |
| `GET` | `/decisions/{id}/conflicts` | Contradiction preview on shared systems |
| `POST` | `/remember` | Submit explicit memory → Kafka → extractor → graph |
| `POST` | `/remember` | Submit explicit memory → Kafka → extractor → graph (shared lazy Kafka producer) |
| `POST` | `/gdpr/erase` | GDPR Right to Erasure — cascade delete subject memory (`admin` / `gdpr_officer` / `legal`); invalidates per-workspace query cache |
| `POST` | `/webhooks/slack`, `/github`, `/jira`, `/linear` | Connector ingress → Kafka |

Expand Down
36 changes: 28 additions & 8 deletions api/remember.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from __future__ import annotations

import os
import threading
from datetime import UTC, datetime

import structlog
Expand All @@ -19,6 +20,11 @@

KAFKA_TOPIC = "cortex.raw.manual.events"

# Module-level lazy singleton — one Producer for the process lifetime.
# Matches connector pattern (SlackKafkaProducer keeps a long-lived client).
_producer_lock = threading.Lock()
_producer_instance: Producer | None = None


class RememberRequest(BaseModel):
workspace_id: str = Field(description="Workspace to store memory in")
Expand All @@ -38,14 +44,28 @@ class RememberResponse(BaseModel):


def _producer() -> Producer:
servers = os.environ.get("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092")
return Producer(
{
"bootstrap.servers": servers,
"acks": "all",
"enable.idempotence": True,
}
)
"""Return the shared Kafka producer, creating it on first use."""
global _producer_instance
if _producer_instance is not None:
return _producer_instance
with _producer_lock:
if _producer_instance is None:
servers = os.environ.get("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092")
_producer_instance = Producer(
{
"bootstrap.servers": servers,
"acks": "all",
"enable.idempotence": True,
}
)
return _producer_instance


def _reset_producer_for_tests() -> None:
"""Clear the singleton between tests (test helper only)."""
global _producer_instance
with _producer_lock:
_producer_instance = None


@router.post(
Expand Down
21 changes: 21 additions & 0 deletions tests/api/test_remember.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,18 @@

from unittest.mock import MagicMock, patch

import pytest
from fastapi.testclient import TestClient

from api.main import app
from api.remember import _producer, _reset_producer_for_tests


@pytest.fixture(autouse=True)
def _clear_producer_singleton() -> None:
_reset_producer_for_tests()
yield
_reset_producer_for_tests()


@patch("api.remember._producer")
Expand Down Expand Up @@ -38,3 +47,15 @@ def test_remember_rejects_short_content() -> None:
json={"workspace_id": "ws-1", "content": "too short"},
)
assert response.status_code == 422


@patch("api.remember.Producer")
def test_producer_is_lazy_singleton(mock_producer_cls: MagicMock) -> None:
"""Same Producer instance is reused across calls; constructed once."""
mock_producer_cls.return_value = MagicMock(name="shared-producer")

first = _producer()
second = _producer()

assert first is second
mock_producer_cls.assert_called_once()
Loading