mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-04-25 11:18:22 +00:00
* 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>
150 lines
4.7 KiB
Python
150 lines
4.7 KiB
Python
"""Assistants compatibility endpoints.
|
|
|
|
Provides LangGraph Platform-compatible assistants API backed by the
|
|
``langgraph.json`` graph registry and ``config.yaml`` agent definitions.
|
|
|
|
This is a minimal stub that satisfies the ``useStream`` React hook's
|
|
initialization requirements (``assistants.search()`` and ``assistants.get()``).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from datetime import UTC, datetime
|
|
from typing import Any
|
|
|
|
from fastapi import APIRouter, HTTPException
|
|
from pydantic import BaseModel, Field
|
|
|
|
logger = logging.getLogger(__name__)
|
|
router = APIRouter(prefix="/api/assistants", tags=["assistants-compat"])
|
|
|
|
|
|
class AssistantResponse(BaseModel):
|
|
assistant_id: str
|
|
graph_id: str
|
|
name: str
|
|
config: dict[str, Any] = Field(default_factory=dict)
|
|
metadata: dict[str, Any] = Field(default_factory=dict)
|
|
description: str | None = None
|
|
created_at: str = ""
|
|
updated_at: str = ""
|
|
version: int = 1
|
|
|
|
|
|
class AssistantSearchRequest(BaseModel):
|
|
graph_id: str | None = None
|
|
name: str | None = None
|
|
metadata: dict[str, Any] | None = None
|
|
limit: int = 10
|
|
offset: int = 0
|
|
|
|
|
|
def _get_default_assistant() -> AssistantResponse:
|
|
"""Return the default lead_agent assistant."""
|
|
now = datetime.now(UTC).isoformat()
|
|
return AssistantResponse(
|
|
assistant_id="lead_agent",
|
|
graph_id="lead_agent",
|
|
name="lead_agent",
|
|
config={},
|
|
metadata={"created_by": "system"},
|
|
description="DeerFlow lead agent",
|
|
created_at=now,
|
|
updated_at=now,
|
|
version=1,
|
|
)
|
|
|
|
|
|
def _list_assistants() -> list[AssistantResponse]:
|
|
"""List all available assistants from config."""
|
|
assistants = [_get_default_assistant()]
|
|
|
|
# Also include custom agents from config.yaml agents directory
|
|
try:
|
|
from deerflow.config.agents_config import list_custom_agents
|
|
|
|
for agent_cfg in list_custom_agents():
|
|
now = datetime.now(UTC).isoformat()
|
|
assistants.append(
|
|
AssistantResponse(
|
|
assistant_id=agent_cfg.name,
|
|
graph_id="lead_agent", # All agents use the same graph
|
|
name=agent_cfg.name,
|
|
config={},
|
|
metadata={"created_by": "user"},
|
|
description=agent_cfg.description or "",
|
|
created_at=now,
|
|
updated_at=now,
|
|
version=1,
|
|
)
|
|
)
|
|
except Exception:
|
|
logger.debug("Could not load custom agents for assistants list")
|
|
|
|
return assistants
|
|
|
|
|
|
@router.post("/search", response_model=list[AssistantResponse])
|
|
async def search_assistants(body: AssistantSearchRequest | None = None) -> list[AssistantResponse]:
|
|
"""Search assistants.
|
|
|
|
Returns all registered assistants (lead_agent + custom agents from config).
|
|
"""
|
|
assistants = _list_assistants()
|
|
|
|
if body and body.graph_id:
|
|
assistants = [a for a in assistants if a.graph_id == body.graph_id]
|
|
if body and body.name:
|
|
assistants = [a for a in assistants if body.name.lower() in a.name.lower()]
|
|
|
|
offset = body.offset if body else 0
|
|
limit = body.limit if body else 10
|
|
return assistants[offset : offset + limit]
|
|
|
|
|
|
@router.get("/{assistant_id}", response_model=AssistantResponse)
|
|
async def get_assistant_compat(assistant_id: str) -> AssistantResponse:
|
|
"""Get an assistant by ID."""
|
|
for a in _list_assistants():
|
|
if a.assistant_id == assistant_id:
|
|
return a
|
|
raise HTTPException(status_code=404, detail=f"Assistant {assistant_id} not found")
|
|
|
|
|
|
@router.get("/{assistant_id}/graph")
|
|
async def get_assistant_graph(assistant_id: str) -> dict:
|
|
"""Get the graph structure for an assistant.
|
|
|
|
Returns a minimal graph description. Full graph introspection is
|
|
not supported in the Gateway — this stub satisfies SDK validation.
|
|
"""
|
|
found = any(a.assistant_id == assistant_id for a in _list_assistants())
|
|
if not found:
|
|
raise HTTPException(status_code=404, detail=f"Assistant {assistant_id} not found")
|
|
|
|
return {
|
|
"graph_id": "lead_agent",
|
|
"nodes": [],
|
|
"edges": [],
|
|
}
|
|
|
|
|
|
@router.get("/{assistant_id}/schemas")
|
|
async def get_assistant_schemas(assistant_id: str) -> dict:
|
|
"""Get JSON schemas for an assistant's input/output/state.
|
|
|
|
Returns empty schemas — full introspection not supported in Gateway.
|
|
"""
|
|
found = any(a.assistant_id == assistant_id for a in _list_assistants())
|
|
if not found:
|
|
raise HTTPException(status_code=404, detail=f"Assistant {assistant_id} not found")
|
|
|
|
return {
|
|
"graph_id": "lead_agent",
|
|
"input_schema": {},
|
|
"output_schema": {},
|
|
"state_schema": {},
|
|
"config_schema": {},
|
|
}
|