deerflow-code/offline-backend-20260512/backend/packages/harness/deerflow/persistence/thread_meta/base.py
2026-09-07 18:24:55 +08:00

145 lines
5.0 KiB
Python

"""Abstract interface for thread metadata storage.
Implementations:
- ThreadMetaRepository: SQL-backed (sqlite / postgres via SQLAlchemy)
- MemoryThreadMetaStore: wraps LangGraph BaseStore (memory mode)
All mutating and querying methods accept a ``user_id`` parameter with
three-state semantics (see :mod:`deerflow.runtime.user_context`):
- ``AUTO`` (default): resolve from the request-scoped contextvar.
- Explicit ``str``: use the provided value verbatim.
- Explicit ``None``: bypass owner filtering (migration/CLI only).
"""
from __future__ import annotations
import abc
from typing import Any
from deerflow.runtime.user_context import AUTO, _AutoSentinel
# Configuration wizards / machine workers that must not appear in the default
# conversation list. Prefer ``metadata.system=true`` on create; this set also
# covers older rows that only set ``thread_type``.
EXCLUDED_SYSTEM_THREAD_TYPES = frozenset(
{
"scheduler",
"roundtable",
"notebook",
"deep_research",
"writing_setup",
"report_structure_setup",
}
)
def is_excluded_system_thread(metadata: dict[str, Any] | None) -> bool:
"""Return True when a thread should be hidden from user conversation lists."""
md = metadata or {}
if md.get("system"):
return True
thread_type = md.get("thread_type")
return isinstance(thread_type, str) and thread_type in EXCLUDED_SYSTEM_THREAD_TYPES
class ThreadMetaStore(abc.ABC):
@abc.abstractmethod
async def create(
self,
thread_id: str,
*,
assistant_id: str | None = None,
user_id: str | None | _AutoSentinel = AUTO,
display_name: str | None = None,
metadata: dict | None = None,
) -> dict:
"""Create one metadata row.
Implementations must assemble the returned dict from the values passed
in, without a post-commit read-back: on MySQL such a read can be routed
to a lagging replica behind a read/write-splitting endpoint and fail.
"""
pass
@abc.abstractmethod
async def get(self, thread_id: str, *, user_id: str | None | _AutoSentinel = AUTO) -> dict | None:
pass
@abc.abstractmethod
async def search(
self,
*,
metadata: dict | None = None,
status: str | None = None,
limit: int = 100,
offset: int = 0,
user_id: str | None | _AutoSentinel = AUTO,
exclude_system: bool = False,
query: str | None = None,
) -> list[dict]:
"""Search threads ordered by ``updated_at`` desc.
``exclude_system=True`` drops machine / wizard threads (``metadata.system``
or known ``thread_type`` values such as scheduler / report_structure_setup)
so user-facing conversation lists are not crowded out.
``query`` is a case-insensitive substring match on the thread title
(``display_name``). ``limit``/``offset`` apply to the *filtered*
result set.
"""
pass
@abc.abstractmethod
async def count(
self,
*,
metadata: dict | None = None,
status: str | None = None,
user_id: str | None | _AutoSentinel = AUTO,
exclude_system: bool = False,
query: str | None = None,
) -> int:
"""Count threads matching the same filters as :meth:`search`.
Backs pagination UIs that need a total alongside per-page results.
"""
pass
@abc.abstractmethod
async def update_display_name(self, thread_id: str, display_name: str, *, user_id: str | None | _AutoSentinel = AUTO) -> None:
pass
@abc.abstractmethod
async def update_status(self, thread_id: str, status: str, *, user_id: str | None | _AutoSentinel = AUTO) -> None:
pass
@abc.abstractmethod
async def update_assistant_id(self, thread_id: str, assistant_id: str, *, user_id: str | None | _AutoSentinel = AUTO) -> None:
"""Backfill ``assistant_id`` for a thread row.
Used by ``start_run`` when the row was pre-created via ``POST /threads``
without an assistant_id (the SDK's ``client.threads.create()`` doesn't
pass it), and the first run on the thread reveals the assistant. No-op
if the thread does not exist or the owner check fails.
"""
pass
@abc.abstractmethod
async def update_metadata(self, thread_id: str, metadata: dict, *, user_id: str | None | _AutoSentinel = AUTO) -> None:
"""Merge ``metadata`` into the thread's metadata field.
Existing keys are overwritten by the new values; keys absent from
``metadata`` are preserved. No-op if the thread does not exist
or the owner check fails.
"""
pass
@abc.abstractmethod
async def check_access(self, thread_id: str, user_id: str, *, require_existing: bool = False) -> bool:
"""Check if ``user_id`` has access to ``thread_id``."""
pass
@abc.abstractmethod
async def delete(self, thread_id: str, *, user_id: str | None | _AutoSentinel = AUTO) -> None:
pass