"""Background auto-sediment queue (phase 3). A single asyncio worker drains a bounded queue of finished conversations and sediments the worthwhile ones into the knowledge base — entirely off the chat request path, so it never adds latency to a reply. Enqueue happens from ``KnowledgeRagMiddleware.after_agent`` when ``knowledge.auto_ingest_enabled``. Value judgment is intentionally cheap (min assistant chars + dedup via ``skip_existing``); the LLM still distills the kept conversations via the normal ``capture_thread`` path. """ from __future__ import annotations import asyncio import contextlib import logging from collections.abc import Callable from dataclasses import dataclass, field from typing import Any logger = logging.getLogger(__name__) _QUEUE: AutoIngestQueue | None = None @dataclass class _QueueUser: """Minimal CurrentUser stand-in (only ``.id`` is needed for context).""" id: str @dataclass class AutoIngestItem: thread_id: str user_id: str messages: list[dict[str, Any]] thread_title: str | None = None reference_batches: list[dict[str, Any]] = field(default_factory=list) class AutoIngestQueue: """Bounded background queue + single worker for auto-sedimentation.""" def __init__(self, service_provider: Callable[[], Any], *, min_chars: int = 300, maxsize: int = 256) -> None: self._provider = service_provider self._min_chars = min_chars self._queue: asyncio.Queue[AutoIngestItem] = asyncio.Queue(maxsize=maxsize) self._task: asyncio.Task | None = None def start(self) -> None: if self._task is None or self._task.done(): self._task = asyncio.create_task(self._worker(), name="knowledge-auto-ingest") logger.info("Knowledge auto-ingest worker started") async def stop(self) -> None: if self._task is not None: self._task.cancel() with contextlib.suppress(asyncio.CancelledError): await self._task self._task = None def enqueue(self, item: AutoIngestItem) -> None: try: self._queue.put_nowait(item) except asyncio.QueueFull: logger.warning("Knowledge auto-ingest queue full; dropping thread %s", item.thread_id) def _worth_ingesting(self, item: AutoIngestItem) -> bool: chars = 0 for m in item.messages: if str(m.get("type") or m.get("role") or "").lower() == "ai": content = m.get("content") chars += len(content) if isinstance(content, str) else len(str(content or "")) return chars >= self._min_chars async def _process(self, item: AutoIngestItem) -> None: service = self._provider() if service is None: return if not self._worth_ingesting(item): logger.debug("Auto-ingest skipped (below min chars) thread=%s", item.thread_id) return from deerflow.runtime.user_context import reset_current_user, set_current_user token = set_current_user(_QueueUser(id=item.user_id)) try: result = await service.capture_thread( thread_id=item.thread_id, messages=item.messages, thread_title=item.thread_title, status="draft", # auto-sedimented notes land as drafts for review include_sources=True, reference_batches=item.reference_batches, created_by=item.user_id, skip_existing=True, ) if result.get("created"): logger.info("Auto-ingested thread %s -> note %s", item.thread_id, result["note"]["id"]) except Exception: logger.exception("Auto-ingest failed for thread %s", item.thread_id) finally: reset_current_user(token) async def _worker(self) -> None: while True: item = await self._queue.get() try: await self._process(item) except asyncio.CancelledError: raise except Exception: logger.exception("Auto-ingest worker error") finally: self._queue.task_done() def set_auto_ingest_queue(queue: AutoIngestQueue | None) -> None: global _QUEUE _QUEUE = queue def get_auto_ingest_queue() -> AutoIngestQueue | None: return _QUEUE