mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-05-20 15:11:09 +00:00
34e835bc33
* feat(gateway): implement LangGraph Platform API in Gateway, replace langgraph-cli Implement all core LangGraph Platform API endpoints in the Gateway, allowing it to fully replace the langgraph-cli dev server for local development. This eliminates a heavyweight dependency and simplifies the development stack. Changes: - Add runs lifecycle endpoints (create, stream, wait, cancel, join) - Add threads CRUD and search endpoints - Add assistants compatibility endpoints (search, get, graph, schemas) - Add StreamBridge (in-memory pub/sub for SSE) and async provider - Add RunManager with atomic create_or_reject (eliminates TOCTOU race) - Add worker with interrupt/rollback cancel actions and runtime context injection - Route /api/langgraph/* to Gateway in nginx config - Skip langgraph-cli startup by default (SKIP_LANGGRAPH_SERVER=0 to restore) - Add unit tests for RunManager, SSE format, and StreamBridge * fix: drain bridge queue on client disconnect to prevent backpressure When on_disconnect=continue, keep consuming events from the bridge without yielding, so the worker is not blocked by a full queue. Only on_disconnect=cancel breaks out immediately. Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * fix: remove pytest import Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * fix: Fix default stream_mode to ["values", "messages-tuple"] Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * fix: Remove unused if_exists field from ThreadCreateRequest Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * fix: address review comments on gateway LangGraph API - Mount runs.py router in app.py (missing include_router) - Normalize interrupt_before/after "*" to node list before run_agent() - Use entry.id for SSE event ID instead of counter - Drain bridge queue on disconnect when on_disconnect=continue - Reuse serialization helper in wait_run() for consistent wire format - Reject unsupported multitask_strategy with 400 - Remove SKIP_LANGGRAPH_SERVER fallback, always use Gateway * feat: extract app.state access into deps.py Encapsulate read/write operations for singleton objects (RunManager, StreamBridge, checkpointer) held in app.state into a shared utility, reducing repeated access patterns across router modules. * feat: extract deerflow.runtime.serialization module with tests Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * refactor: replace duplicated serialization with deerflow.runtime.serialization Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * feat: extract app/gateway/services.py with run lifecycle logic Create a service layer that centralizes SSE formatting, input/config normalization, and run lifecycle management. Router modules will delegate to these functions instead of using private cross-imported helpers. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * refactor: wire routers to use services layer, remove cross-module private imports Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * style: apply ruff formatting to refactored files Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * feat(runtime): support LangGraph dev server and add compat route - Enable official LangGraph dev server for local development workflow - Decouple runtime components from agents package for better separation - Provide gateway-backed fallback route when dev server is skipped - Simplify lifecycle management using context manager in gateway * feat(runtime): add Store providers with auto-backend selection - Add async_provider.py and provider.py under deerflow/runtime/store/ - Support memory, sqlite, postgres backends matching checkpointer config - Integrate into FastAPI lifespan via AsyncExitStack in deps.py - Replace hardcoded InMemoryStore with config-driven factory * refactor(gateway): migrate thread management from checkpointer to Store and resolve multiple endpoint failures - Add Store-backed CRUD helpers (_store_get, _store_put, _store_upsert) - Replace checkpoint-scanning search with two-phase strategy: phase 1 reads Store (O(threads)), phase 2 backfills from checkpointer for legacy/LangGraph Server threads with lazy migration - Extend Store record schema with values field for title persistence - Sync thread title from checkpoint to Store after run completion - Fix /threads/{id}/runs/{run_id}/stream 405 by accepting both GET and POST methods; POST handles interrupt/rollback actions - Fix /threads/{id}/state 500 by separating read_config and write_config, adding checkpoint_ns to configurable, and shallow-copying checkpoint/metadata before mutation - Sync title to Store on state update for immediate search reflection - Move _upsert_thread_in_store into services.py, remove duplicate logic - Add _sync_thread_title_after_run: await run task, read final checkpoint title, write back to Store record - Spawn title sync as background task from start_run when Store exists * refactor(runtime): deduplicate store and checkpointer provider logic Extract _ensure_sqlite_parent_dir() helper into checkpointer/provider.py and use it in all three places that previously inlined the same mkdir logic. Consolidate duplicate error constants in store/async_provider.py by importing from store/provider.py instead of redefining them. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * refactor(runtime): move SQLite helpers to runtime/store, checkpointer imports from store _resolve_sqlite_conn_str and _ensure_sqlite_parent_dir now live in runtime/store/provider.py. agents/checkpointer/provider and agents/checkpointer/async_provider import from there, reversing the previous dependency direction (store → checkpointer becomes checkpointer → store). Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * refactor(runtime): extract SQLite helpers into runtime/store/_sqlite_utils.py Move resolve_sqlite_conn_str and ensure_sqlite_parent_dir out of checkpointer/provider.py into a dedicated _sqlite_utils module. Functions are now public (no underscore prefix), making cross-module imports semantically correct. All four provider files import from the single shared location. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(gateway): use adelete_thread to fully remove thread checkpoints on delete AsyncSqliteSaver has no adelete method — the previous hasattr check always evaluated to False, silently leaving all checkpoint rows in the database. Switch to adelete_thread(thread_id) which deletes every checkpoint and pending-write row for the thread across all namespaces (including sub-graph checkpoints). Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix(gateway): remove dead bridge_cm/ckpt_cm code and fix StrEnum lint app.py had unreachable code after the async-with lifespan refactor: bridge_cm and ckpt_cm were referenced but never defined (F821), and the channel service startup/shutdown was outside the langgraph_runtime block so it never ran. Move channel service lifecycle inside the async-with block where it belongs. Replace str+Enum inheritance in RunStatus and DisconnectMode with StrEnum as suggested by UP042. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * style: format with ruff --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> Co-authored-by: JeffJiang <for-eleven@hotmail.com> Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com> Co-authored-by: Willem Jiang <willem.jiang@gmail.com>
189 lines
6.7 KiB
Python
189 lines
6.7 KiB
Python
"""Sync Store factory.
|
|
|
|
Provides a **sync singleton** and a **sync context manager** for CLI tools
|
|
and the embedded :class:`~deerflow.client.DeerFlowClient`.
|
|
|
|
The backend mirrors the configured checkpointer so that both always use the
|
|
same persistence technology. Supported backends: memory, sqlite, postgres.
|
|
|
|
Usage::
|
|
|
|
from deerflow.runtime.store.provider import get_store, store_context
|
|
|
|
# Singleton — reused across calls, closed on process exit
|
|
store = get_store()
|
|
|
|
# One-shot — fresh connection, closed on block exit
|
|
with store_context() as store:
|
|
store.put(("ns",), "key", {"value": 1})
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import contextlib
|
|
import logging
|
|
from collections.abc import Iterator
|
|
|
|
from langgraph.store.base import BaseStore
|
|
|
|
from deerflow.config.app_config import get_app_config
|
|
from deerflow.runtime.store._sqlite_utils import ensure_sqlite_parent_dir, resolve_sqlite_conn_str
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Error message constants
|
|
# ---------------------------------------------------------------------------
|
|
|
|
SQLITE_STORE_INSTALL = "langgraph-checkpoint-sqlite is required for the SQLite store. Install it with: uv add langgraph-checkpoint-sqlite"
|
|
POSTGRES_STORE_INSTALL = "langgraph-checkpoint-postgres is required for the PostgreSQL store. Install it with: uv add langgraph-checkpoint-postgres psycopg[binary] psycopg-pool"
|
|
POSTGRES_CONN_REQUIRED = "checkpointer.connection_string is required for the postgres backend"
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Sync factory
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@contextlib.contextmanager
|
|
def _sync_store_cm(config) -> Iterator[BaseStore]:
|
|
"""Context manager that creates and tears down a sync Store.
|
|
|
|
The ``config`` argument is a
|
|
:class:`~deerflow.config.checkpointer_config.CheckpointerConfig` instance —
|
|
the same object used by the checkpointer factory.
|
|
"""
|
|
if config.type == "memory":
|
|
from langgraph.store.memory import InMemoryStore
|
|
|
|
logger.info("Store: using InMemoryStore (in-process, not persistent)")
|
|
yield InMemoryStore()
|
|
return
|
|
|
|
if config.type == "sqlite":
|
|
try:
|
|
from langgraph.store.sqlite import SqliteStore
|
|
except ImportError as exc:
|
|
raise ImportError(SQLITE_STORE_INSTALL) from exc
|
|
|
|
conn_str = resolve_sqlite_conn_str(config.connection_string or "store.db")
|
|
ensure_sqlite_parent_dir(conn_str)
|
|
|
|
with SqliteStore.from_conn_string(conn_str) as store:
|
|
store.setup()
|
|
logger.info("Store: using SqliteStore (%s)", conn_str)
|
|
yield store
|
|
return
|
|
|
|
if config.type == "postgres":
|
|
try:
|
|
from langgraph.store.postgres import PostgresStore # type: ignore[import]
|
|
except ImportError as exc:
|
|
raise ImportError(POSTGRES_STORE_INSTALL) from exc
|
|
|
|
if not config.connection_string:
|
|
raise ValueError(POSTGRES_CONN_REQUIRED)
|
|
|
|
with PostgresStore.from_conn_string(config.connection_string) as store:
|
|
store.setup()
|
|
logger.info("Store: using PostgresStore")
|
|
yield store
|
|
return
|
|
|
|
raise ValueError(f"Unknown store backend type: {config.type!r}")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Sync singleton
|
|
# ---------------------------------------------------------------------------
|
|
|
|
_store: BaseStore | None = None
|
|
_store_ctx = None # open context manager keeping the connection alive
|
|
|
|
|
|
def get_store() -> BaseStore:
|
|
"""Return the global sync Store singleton, creating it on first call.
|
|
|
|
Returns an :class:`~langgraph.store.memory.InMemoryStore` when no
|
|
checkpointer is configured in *config.yaml* (emits a WARNING in that case).
|
|
|
|
Raises:
|
|
ImportError: If the required package for the configured backend is not installed.
|
|
ValueError: If ``connection_string`` is missing for a backend that requires it.
|
|
"""
|
|
global _store, _store_ctx
|
|
|
|
if _store is not None:
|
|
return _store
|
|
|
|
# Lazily load app config, mirroring the checkpointer singleton pattern so
|
|
# that tests that set the global checkpointer config explicitly remain isolated.
|
|
from deerflow.config.app_config import _app_config
|
|
from deerflow.config.checkpointer_config import get_checkpointer_config
|
|
|
|
config = get_checkpointer_config()
|
|
|
|
if config is None and _app_config is None:
|
|
try:
|
|
get_app_config()
|
|
except FileNotFoundError:
|
|
pass
|
|
config = get_checkpointer_config()
|
|
|
|
if config is None:
|
|
from langgraph.store.memory import InMemoryStore
|
|
|
|
logger.warning("No 'checkpointer' section in config.yaml — using InMemoryStore for the store. Thread list will be lost on server restart. Configure a sqlite or postgres backend for persistence.")
|
|
_store = InMemoryStore()
|
|
return _store
|
|
|
|
_store_ctx = _sync_store_cm(config)
|
|
_store = _store_ctx.__enter__()
|
|
return _store
|
|
|
|
|
|
def reset_store() -> None:
|
|
"""Reset the sync singleton, forcing recreation on the next call.
|
|
|
|
Closes any open backend connections and clears the cached instance.
|
|
Useful in tests or after a configuration change.
|
|
"""
|
|
global _store, _store_ctx
|
|
if _store_ctx is not None:
|
|
try:
|
|
_store_ctx.__exit__(None, None, None)
|
|
except Exception:
|
|
logger.warning("Error during store cleanup", exc_info=True)
|
|
_store_ctx = None
|
|
_store = None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Sync context manager
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@contextlib.contextmanager
|
|
def store_context() -> Iterator[BaseStore]:
|
|
"""Sync context manager that yields a Store and cleans up on exit.
|
|
|
|
Unlike :func:`get_store`, this does **not** cache the instance — each
|
|
``with`` block creates and destroys its own connection. Use it in CLI
|
|
scripts or tests where you want deterministic cleanup::
|
|
|
|
with store_context() as store:
|
|
store.put(("threads",), thread_id, {...})
|
|
|
|
Yields an :class:`~langgraph.store.memory.InMemoryStore` when no
|
|
checkpointer is configured in *config.yaml*.
|
|
"""
|
|
config = get_app_config()
|
|
if config.checkpointer is None:
|
|
from langgraph.store.memory import InMemoryStore
|
|
|
|
logger.warning("No 'checkpointer' section in config.yaml — using InMemoryStore for the store. Thread list will be lost on server restart. Configure a sqlite or postgres backend for persistence.")
|
|
yield InMemoryStore()
|
|
return
|
|
|
|
with _sync_store_cm(config.checkpointer) as store:
|
|
yield store
|