From fd03b244c8c1ab53f2d225714a6f4a8e780c9f8d Mon Sep 17 00:00:00 2001 From: Abhinaysai Kamineni <66816045+askmy-stack@users.noreply.github.com> Date: Thu, 16 Jul 2026 07:55:32 -0400 Subject: [PATCH] Reuse a single Kafka producer for POST /remember Avoid constructing a new confluent_kafka Producer on every request; lazy-init a process-wide singleton with double-checked locking. Co-authored-by: Cursor --- README.md | 2 +- api/remember.py | 36 ++++++++++++++++++++++++++++-------- tests/api/test_remember.py | 21 +++++++++++++++++++++ 3 files changed, 50 insertions(+), 9 deletions(-) diff --git a/README.md b/README.md index 834f468..3ccfbf9 100644 --- a/README.md +++ b/README.md @@ -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 | diff --git a/api/remember.py b/api/remember.py index 9b957e0..102d253 100644 --- a/api/remember.py +++ b/api/remember.py @@ -3,6 +3,7 @@ from __future__ import annotations import os +import threading from datetime import UTC, datetime import structlog @@ -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") @@ -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( diff --git a/tests/api/test_remember.py b/tests/api/test_remember.py index 17b97ca..8ac8e2f 100644 --- a/tests/api/test_remember.py +++ b/tests/api/test_remember.py @@ -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") @@ -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()