126 lines
4.3 KiB
Python
126 lines
4.3 KiB
Python
"""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
|