1845 lines
70 KiB
Python
1845 lines
70 KiB
Python
import asyncio
|
||
import logging
|
||
import os
|
||
import queue
|
||
import sys
|
||
import threading
|
||
from collections.abc import AsyncGenerator
|
||
from contextlib import asynccontextmanager
|
||
from datetime import datetime
|
||
from time import perf_counter
|
||
|
||
from dotenv import load_dotenv
|
||
from fastapi import FastAPI
|
||
from fastapi.middleware.cors import CORSMiddleware
|
||
|
||
# Load .env BEFORE any module that reads os.environ at import time.
|
||
# Without this, GATEWAY_CORS_ORIGINS / AUTH_JWT_SECRET / DEER_FLOW_AUTH_DISABLED
|
||
# placed in /workspace/.env are silently ignored when the container is started
|
||
# without explicit `-e` flags, leading to ephemeral JWT secrets, missing CORS
|
||
# origins, and "Failed to fetch" in the browser.
|
||
load_dotenv()
|
||
|
||
from app.gateway.auth.disabled_mode import is_auth_disabled # noqa: E402
|
||
from app.gateway.auth_middleware import AuthMiddleware
|
||
from app.gateway.config import get_gateway_config
|
||
from app.gateway.csrf_middleware import CSRFMiddleware
|
||
from app.gateway.deps import langgraph_runtime
|
||
from app.gateway.public_knowledge_cors import PublicKnowledgeCORSMiddleware
|
||
from app.gateway.routers import (
|
||
admin_active_runs,
|
||
admin_users,
|
||
agents,
|
||
ai_writing,
|
||
artifacts,
|
||
assistant_knowledge,
|
||
assistants_compat,
|
||
auth,
|
||
browser_context,
|
||
business_mapping,
|
||
channels,
|
||
checkpoint_migration,
|
||
compact_menu_overrides,
|
||
dashboard_sessions,
|
||
deep_research,
|
||
external_llmwiki,
|
||
feedback,
|
||
fixed_questions,
|
||
flow_extract,
|
||
html_page_favorites,
|
||
intent,
|
||
knowledge,
|
||
knowledge_ingest,
|
||
light_apps,
|
||
llm_metrics,
|
||
llmwiki,
|
||
llmwiki_index,
|
||
mcp,
|
||
memory,
|
||
menu_overrides,
|
||
model_management,
|
||
models,
|
||
multi_agent,
|
||
notifications,
|
||
open_3qfx,
|
||
open_chat,
|
||
parallel_agent_tasks,
|
||
parallel_agents,
|
||
position_roles,
|
||
position_roundtable,
|
||
positions,
|
||
public_config,
|
||
public_embed,
|
||
public_knowledge_vector_search,
|
||
public_leaderboard,
|
||
public_scheduled_tasks,
|
||
recommend,
|
||
recommended_questions,
|
||
report_collaboration,
|
||
report_structures,
|
||
roundtable_artifact_submissions,
|
||
roundtable_chains,
|
||
roundtable_diagnostics,
|
||
roundtable_draft_shares,
|
||
roundtable_drafts,
|
||
roundtable_jobs,
|
||
roundtable_task_drafts,
|
||
runs,
|
||
scheduled_tasks,
|
||
sensitive_words,
|
||
sentiment_agent,
|
||
skill_addresses,
|
||
skill_knowledge,
|
||
skills,
|
||
suggestions,
|
||
system_settings,
|
||
tags,
|
||
task_buttons,
|
||
task_reports,
|
||
taskcop_tasks,
|
||
thread_runs,
|
||
thread_shares,
|
||
threads,
|
||
tool_metrics,
|
||
uploads,
|
||
user_preferences,
|
||
user_prompts,
|
||
workflow_data_sources,
|
||
workflow_embed,
|
||
workflow_planning,
|
||
workflow_resources,
|
||
workflow_runs,
|
||
workflow_studio,
|
||
workflows,
|
||
workflows_coze_compat,
|
||
writing,
|
||
)
|
||
from deerflow.config import app_config as deerflow_app_config
|
||
from deerflow.config.app_config import apply_logging_level
|
||
from deerflow.runtime.scheduler import ScheduledTaskService, set_scheduled_task_service
|
||
|
||
AppConfig = deerflow_app_config.AppConfig
|
||
get_app_config = deerflow_app_config.get_app_config
|
||
|
||
# Default logging; lifespan overrides from config.yaml log_level.
|
||
logging.basicConfig(
|
||
level=logging.INFO,
|
||
format="%(asctime)s - %(name)s - %(levelname)s - %(message)s",
|
||
datefmt="%Y-%m-%d %H:%M:%S",
|
||
)
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Upper bound (seconds) each lifespan shutdown hook is allowed to run.
|
||
# Bounds worker exit time so uvicorn's reload supervisor does not keep
|
||
# firing signals into a worker that is stuck waiting for shutdown cleanup.
|
||
_SHUTDOWN_HOOK_TIMEOUT_SECONDS = 5.0
|
||
|
||
|
||
async def _ensure_admin_user(app: FastAPI) -> None:
|
||
"""Startup hook: handle first boot and migrate orphan threads otherwise.
|
||
|
||
After admin creation, migrate orphan threads from the LangGraph
|
||
store (metadata.user_id unset) to the admin account. This is the
|
||
"no-auth ? with-auth" upgrade path: users who ran DeerFlow without
|
||
authentication have existing LangGraph thread data that needs an
|
||
owner assigned.
|
||
First boot (no admin exists):
|
||
- Does NOT create any user accounts automatically.
|
||
- The operator must visit ``/setup`` to create the first admin.
|
||
|
||
Subsequent boots (admin already exists):
|
||
- Runs the one-time "no-auth ? with-auth" orphan thread migration for
|
||
existing LangGraph thread metadata that has no owner_id.
|
||
|
||
No SQL persistence migration is needed: the four user_id columns
|
||
(threads_meta, runs, run_events, feedback) only come into existence
|
||
alongside the auth module via create_all, so freshly created tables
|
||
never contain NULL-owner rows.
|
||
"""
|
||
from sqlalchemy import select
|
||
|
||
from app.gateway.deps import get_local_provider
|
||
from deerflow.persistence.engine import get_session_factory
|
||
from deerflow.persistence.user.model import UserRow
|
||
|
||
try:
|
||
provider = get_local_provider()
|
||
except RuntimeError:
|
||
# Auth persistence may not be initialized in some test/boot paths.
|
||
# Skip admin migration work rather than failing gateway startup.
|
||
logger.warning("Auth persistence not ready; skipping admin bootstrap check")
|
||
return
|
||
|
||
sf = get_session_factory()
|
||
if sf is None:
|
||
return
|
||
|
||
admin_count = await provider.count_admin_users()
|
||
|
||
if admin_count == 0:
|
||
logger.info("=" * 60)
|
||
logger.info(" First boot detected ? no admin account exists.")
|
||
logger.info(" Visit /setup to complete admin account creation.")
|
||
logger.info("=" * 60)
|
||
return
|
||
|
||
# Admin already exists ? run orphan thread migration for any
|
||
# LangGraph thread metadata that pre-dates the auth module.
|
||
async with sf() as session:
|
||
stmt = select(UserRow).where(UserRow.system_role == "admin").limit(1)
|
||
row = (await session.execute(stmt)).scalar_one_or_none()
|
||
|
||
if row is None:
|
||
return # Should not happen (admin_count > 0 above), but be safe.
|
||
|
||
admin_id = str(row.id)
|
||
|
||
# LangGraph store orphan migration ? non-fatal.
|
||
# This covers the "no-auth ? with-auth" upgrade path for users
|
||
# whose existing LangGraph thread metadata has no user_id set.
|
||
store = getattr(app.state, "store", None)
|
||
if store is not None:
|
||
try:
|
||
migrated = await _migrate_orphaned_threads(store, admin_id)
|
||
if migrated:
|
||
logger.info("Migrated %d orphan LangGraph thread(s) to admin", migrated)
|
||
except Exception:
|
||
logger.exception("LangGraph thread migration failed (non-fatal)")
|
||
|
||
|
||
async def _iter_store_items(store, namespace, *, page_size: int = 500):
|
||
"""Paginated async iterator over a LangGraph store namespace.
|
||
|
||
Replaces the old hardcoded ``limit=1000`` call with a cursor-style
|
||
loop so that environments with more than one page of orphans do
|
||
not silently lose data. Terminates when a page is empty OR when a
|
||
short page arrives (indicating the last page).
|
||
"""
|
||
offset = 0
|
||
while True:
|
||
batch = await store.asearch(namespace, limit=page_size, offset=offset)
|
||
if not batch:
|
||
return
|
||
for item in batch:
|
||
yield item
|
||
if len(batch) < page_size:
|
||
return
|
||
offset += page_size
|
||
|
||
|
||
async def _sync_legacy_skills(app: FastAPI) -> None:
|
||
"""Idempotently register existing on-disk ``skills/custom/*`` directories
|
||
in the ``skills`` DB table, owned by the first admin and published=True.
|
||
|
||
Without this step, after upgrading to per-user skill ownership the existing
|
||
custom skill folders would not be visible to anyone (the SkillStore
|
||
visibility filter relies on a DB row). Running this once on every boot is
|
||
cheap (folder scan + one ``ensure_legacy`` per name) and idempotent.
|
||
|
||
Non-fatal: a failure here MUST NOT block server startup.
|
||
"""
|
||
from sqlalchemy import select
|
||
|
||
from deerflow.persistence.engine import get_session_factory
|
||
from deerflow.persistence.user.model import UserRow
|
||
from deerflow.skills.storage import get_or_new_skill_storage
|
||
from deerflow.skills.types import SkillCategory
|
||
|
||
sf = get_session_factory()
|
||
if sf is None:
|
||
return
|
||
|
||
skill_store = getattr(app.state, "skill_store", None)
|
||
if skill_store is None:
|
||
return
|
||
|
||
async with sf() as session:
|
||
stmt = select(UserRow).where(UserRow.system_role == "admin").limit(1)
|
||
row = (await session.execute(stmt)).scalar_one_or_none()
|
||
if row is None:
|
||
# No admin yet ? leave ownership null; first admin is responsible for
|
||
# claiming via the UI once they exist.
|
||
admin_id: str | None = None
|
||
else:
|
||
admin_id = str(row.id)
|
||
|
||
config = getattr(app.state, "config", None)
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
try:
|
||
skills = storage.load_skills(enabled_only=False)
|
||
except Exception:
|
||
logger.exception("_sync_legacy_skills: could not list on-disk skills")
|
||
return
|
||
|
||
inserted = 0
|
||
for skill in skills:
|
||
if skill.category != SkillCategory.CUSTOM:
|
||
continue
|
||
try:
|
||
existing = await skill_store.get_any(skill.name)
|
||
if existing is not None:
|
||
continue
|
||
await skill_store.ensure_legacy(
|
||
{"name": skill.name, "owner_user_id": admin_id, "published": True}
|
||
)
|
||
inserted += 1
|
||
except Exception:
|
||
logger.warning("_sync_legacy_skills: could not register %s", skill.name, exc_info=True)
|
||
if inserted:
|
||
logger.info("Skill ownership migration: registered %d legacy custom skill(s)", inserted)
|
||
|
||
|
||
async def _migrate_orphaned_threads(store, admin_user_id: str) -> int:
|
||
"""Migrate LangGraph store threads with no user_id to the given admin.
|
||
|
||
Uses cursor pagination so all orphans are migrated regardless of
|
||
count. Returns the number of rows migrated.
|
||
"""
|
||
migrated = 0
|
||
async for item in _iter_store_items(store, ("threads",)):
|
||
metadata = item.value.get("metadata", {})
|
||
if not metadata.get("user_id"):
|
||
metadata["user_id"] = admin_user_id
|
||
item.value["metadata"] = metadata
|
||
await store.aput(("threads",), item.key, item.value)
|
||
migrated += 1
|
||
return migrated
|
||
|
||
|
||
async def _run_db_migrations() -> None:
|
||
"""Run Alembic migrations at startup so schema changes are applied automatically."""
|
||
from pathlib import Path
|
||
|
||
from alembic import command
|
||
from alembic.config import Config
|
||
|
||
import deerflow.persistence as _persistence_pkg
|
||
from deerflow.persistence.engine import get_engine
|
||
|
||
engine = get_engine()
|
||
if engine is None:
|
||
logger.info("Persistence backend=memory ? skipping Alembic migrations")
|
||
return
|
||
|
||
if os.getenv("DEER_FLOW_SKIP_ALEMBIC", "").lower() in {"1", "true", "yes"}:
|
||
logger.warning("DEER_FLOW_SKIP_ALEMBIC is set; skipping Alembic migrations")
|
||
return
|
||
|
||
db_url = engine.url.render_as_string(hide_password=False).replace("%", "%%")
|
||
migrations_dir = str(Path(_persistence_pkg.__file__).parent / "migrations")
|
||
result_queue: queue.Queue[BaseException | None] = queue.Queue(maxsize=1)
|
||
|
||
def _run_sync() -> None:
|
||
try:
|
||
cfg = Config()
|
||
cfg.set_main_option("script_location", migrations_dir)
|
||
cfg.set_main_option("sqlalchemy.url", db_url)
|
||
command.upgrade(cfg, "head")
|
||
except BaseException as exc:
|
||
result_queue.put(exc)
|
||
else:
|
||
result_queue.put(None)
|
||
|
||
try:
|
||
try:
|
||
timeout_s = float(os.getenv("DEER_FLOW_ALEMBIC_TIMEOUT", "60"))
|
||
except ValueError:
|
||
timeout_s = 60.0
|
||
|
||
worker = threading.Thread(target=_run_sync, name="deerflow-alembic-upgrade", daemon=True)
|
||
worker.start()
|
||
try:
|
||
deadline = asyncio.get_running_loop().time() + timeout_s
|
||
while result_queue.empty():
|
||
if asyncio.get_running_loop().time() >= deadline:
|
||
raise TimeoutError
|
||
await asyncio.sleep(0.2)
|
||
result = result_queue.get_nowait()
|
||
except TimeoutError:
|
||
logger.error(
|
||
"Database migration did not finish within %.0fs; continuing startup on create_all schema. "
|
||
"Likely cause: MySQL metadata-lock wait or a slow DDL. Set DEER_FLOW_SKIP_ALEMBIC=1 "
|
||
"to skip startup migrations, or inspect MySQL with SHOW FULL PROCESSLIST.",
|
||
timeout_s,
|
||
)
|
||
return
|
||
if result is not None:
|
||
raise result
|
||
logger.info("Database migrations applied (head)")
|
||
except Exception:
|
||
logger.exception("Database migration failed ? schema may be out of date")
|
||
|
||
|
||
# Skill curator ?????
|
||
# - ????:?????????"??????";??????? should_run_now() ???
|
||
# - ????:?????????????,???????
|
||
_CURATOR_CHECK_INTERVAL_SECONDS = 3600 # 1 ??
|
||
_CURATOR_STARTUP_DELAY_SECONDS = 300 # 5 ??
|
||
|
||
|
||
def _acquire_singleton_lock(lock_path):
|
||
"""???????????,??? worker ?? leader ???
|
||
|
||
?? -> ??????????(?????????,??????/????);
|
||
????(???? worker ??)????? -> ?? None?
|
||
??????? OS ????????
|
||
"""
|
||
try:
|
||
lock_path.parent.mkdir(parents=True, exist_ok=True)
|
||
handle = open(lock_path, "a+")
|
||
except OSError:
|
||
logger.exception("Curator scheduler: ???????")
|
||
return None
|
||
try:
|
||
if sys.platform == "win32":
|
||
import msvcrt
|
||
|
||
handle.seek(0)
|
||
msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1)
|
||
else:
|
||
import fcntl
|
||
|
||
fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||
except OSError:
|
||
handle.close()
|
||
return None
|
||
return handle
|
||
|
||
|
||
def _release_singleton_lock(handle) -> None:
|
||
"""????? _acquire_singleton_lock ?????????"""
|
||
try:
|
||
if sys.platform == "win32":
|
||
import msvcrt
|
||
|
||
handle.seek(0)
|
||
msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1)
|
||
else:
|
||
import fcntl
|
||
|
||
fcntl.flock(handle.fileno(), fcntl.LOCK_UN)
|
||
except OSError:
|
||
pass
|
||
finally:
|
||
try:
|
||
handle.close()
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
async def _curator_scheduler_loop() -> None:
|
||
"""??????:?? _CURATOR_CHECK_INTERVAL_SECONDS ???????????
|
||
|
||
? worker ???????? leader ?? ?? ???????? worker ????
|
||
????,?? worker ????,?? curator ?????????? LLM?
|
||
? cancel ? CancelledError ??? ``except Exception``(?? BaseException
|
||
??,?????),?????????? finally ?????
|
||
"""
|
||
from deerflow.config.runtime_paths import runtime_home
|
||
from deerflow.skills.curator import maybe_run_curator
|
||
|
||
lock_handle = _acquire_singleton_lock(runtime_home() / ".curator_scheduler.lock")
|
||
if lock_handle is None:
|
||
logger.info("Curator scheduler: ??? worker ??????,? worker ??")
|
||
return
|
||
|
||
logger.info("Curator scheduler: ?????,? worker ?? curator ??")
|
||
try:
|
||
# ?????????????,??????
|
||
await asyncio.sleep(_CURATOR_STARTUP_DELAY_SECONDS)
|
||
while True:
|
||
try:
|
||
await maybe_run_curator()
|
||
except Exception:
|
||
logger.exception("Curator scheduler tick failed (non-fatal)")
|
||
await asyncio.sleep(_CURATOR_CHECK_INTERVAL_SECONDS)
|
||
finally:
|
||
_release_singleton_lock(lock_handle)
|
||
|
||
|
||
def _env_bool(name: str, default: bool) -> bool:
|
||
raw = os.getenv(name)
|
||
if raw is None:
|
||
return default
|
||
return raw.strip().lower() not in {"0", "false", "no", "off"}
|
||
|
||
|
||
def _env_int(name: str, default: int, *, alias: str | None = None, minimum: int = 1) -> int:
|
||
raw = os.getenv(name)
|
||
if raw is None and alias:
|
||
raw = os.getenv(alias)
|
||
if raw is None:
|
||
raw = str(default)
|
||
try:
|
||
value = int(raw)
|
||
except ValueError:
|
||
return default
|
||
return max(minimum, value)
|
||
|
||
|
||
def _env_float(name: str, default: float, *, alias: str | None = None, minimum: float = 0.0) -> float:
|
||
raw = os.getenv(name)
|
||
if raw is None and alias:
|
||
raw = os.getenv(alias)
|
||
if raw is None:
|
||
raw = str(default)
|
||
try:
|
||
value = float(raw)
|
||
except ValueError:
|
||
return default
|
||
return max(minimum, value)
|
||
|
||
|
||
_LEADERBOARD_SNAPSHOT_BACKGROUND_ENABLED = _env_bool("ADMIN_LEADERBOARD_BACKGROUND_ENABLED", True)
|
||
_LEADERBOARD_SNAPSHOT_PREWARM_DAYS = _env_int(
|
||
"ADMIN_LEADERBOARD_BACKGROUND_PREWARM_DAYS",
|
||
30,
|
||
)
|
||
_LEADERBOARD_SNAPSHOT_BATCH_SIZE = _env_int(
|
||
"ADMIN_LEADERBOARD_BACKGROUND_BATCH_SIZE",
|
||
1,
|
||
)
|
||
_LEADERBOARD_SNAPSHOT_INTERVAL_SECONDS = _env_float(
|
||
"ADMIN_LEADERBOARD_BACKGROUND_INTERVAL_SECONDS",
|
||
300.0,
|
||
)
|
||
_LEADERBOARD_SNAPSHOT_SLEEP_SECONDS = _env_float(
|
||
"ADMIN_LEADERBOARD_BACKGROUND_SLEEP_SECONDS",
|
||
2.0,
|
||
)
|
||
_LEADERBOARD_SNAPSHOT_RUNNING_TTL_SECONDS = _env_int(
|
||
"ADMIN_LEADERBOARD_BACKGROUND_RUNNING_TTL_SECONDS",
|
||
1800,
|
||
)
|
||
_LEADERBOARD_SNAPSHOT_STARTUP_DELAY_SECONDS = _env_float(
|
||
"ADMIN_LEADERBOARD_BACKGROUND_STARTUP_DELAY_SECONDS",
|
||
60.0,
|
||
)
|
||
# Concurrency (system-pressure) sampler: how often to snapshot the in-flight run
|
||
# count, and how long to keep samples (drives the admin 并发量 line chart, which
|
||
# aggregates per minute/hour). Sampling is in-memory cheap (one tiny INSERT).
|
||
_CONCURRENCY_SAMPLE_ENABLED = _env_bool("CONCURRENCY_SAMPLE_ENABLED", True)
|
||
_CONCURRENCY_SAMPLE_INTERVAL_SECONDS = _env_float(
|
||
"CONCURRENCY_SAMPLE_INTERVAL_SECONDS",
|
||
15.0,
|
||
)
|
||
_CONCURRENCY_SAMPLE_RETENTION_HOURS = _env_float(
|
||
"CONCURRENCY_SAMPLE_RETENTION_HOURS",
|
||
168.0, # 7 days
|
||
)
|
||
# Hour (local Beijing time, 0-23) after which the scheduler force-recomputes the
|
||
# *previous* day's snapshot once, so records that landed near midnight are fully
|
||
# captured. Default 1 ? "???? 1 ????????".
|
||
_LEADERBOARD_FINALIZE_HOUR = min(
|
||
23,
|
||
_env_int("ADMIN_LEADERBOARD_FINALIZE_HOUR", 1, minimum=0),
|
||
)
|
||
|
||
|
||
def _leaderboard_scheduler_sql_for_log(stmt) -> str:
|
||
try:
|
||
return str(stmt.compile(compile_kwargs={"literal_binds": True}))
|
||
except Exception:
|
||
return str(stmt)
|
||
|
||
|
||
async def _execute_leaderboard_scheduler_query(session, step: str, stmt, **context):
|
||
sql = _leaderboard_scheduler_sql_for_log(stmt)
|
||
logger.info(
|
||
"Leaderboard snapshot DB query start: step=%s stat_date=%s job=%s settings_hash=%s started_at=%s sql=%s",
|
||
step,
|
||
context.get("stat_date") or context.get("stat_dates") or "",
|
||
context.get("job_key") or "",
|
||
context.get("settings_hash") or "",
|
||
datetime.now().isoformat(),
|
||
sql,
|
||
)
|
||
start = perf_counter()
|
||
try:
|
||
result = await session.execute(stmt)
|
||
except Exception:
|
||
logger.exception(
|
||
"Leaderboard snapshot DB query failed: step=%s stat_date=%s job=%s settings_hash=%s elapsed_ms=%.2f",
|
||
step,
|
||
context.get("stat_date") or context.get("stat_dates") or "",
|
||
context.get("job_key") or "",
|
||
context.get("settings_hash") or "",
|
||
(perf_counter() - start) * 1000,
|
||
)
|
||
raise
|
||
logger.info(
|
||
"Leaderboard snapshot DB query done: step=%s stat_date=%s job=%s settings_hash=%s finished_at=%s elapsed_ms=%.2f",
|
||
step,
|
||
context.get("stat_date") or context.get("stat_dates") or "",
|
||
context.get("job_key") or "",
|
||
context.get("settings_hash") or "",
|
||
datetime.now().isoformat(),
|
||
(perf_counter() - start) * 1000,
|
||
)
|
||
return result
|
||
|
||
|
||
async def _time_leaderboard_scheduler_db_operation(step: str, operation, *, sql: str = "", **context):
|
||
logger.info(
|
||
"Leaderboard snapshot DB query start: step=%s stat_date=%s job=%s settings_hash=%s started_at=%s sql=%s",
|
||
step,
|
||
context.get("stat_date") or context.get("stat_dates") or "",
|
||
context.get("job_key") or "",
|
||
context.get("settings_hash") or "",
|
||
datetime.now().isoformat(),
|
||
sql,
|
||
)
|
||
start = perf_counter()
|
||
try:
|
||
result = await operation()
|
||
except Exception:
|
||
logger.exception(
|
||
"Leaderboard snapshot DB query failed: step=%s stat_date=%s job=%s settings_hash=%s elapsed_ms=%.2f",
|
||
step,
|
||
context.get("stat_date") or context.get("stat_dates") or "",
|
||
context.get("job_key") or "",
|
||
context.get("settings_hash") or "",
|
||
(perf_counter() - start) * 1000,
|
||
)
|
||
raise
|
||
logger.info(
|
||
"Leaderboard snapshot DB query done: step=%s stat_date=%s job=%s settings_hash=%s finished_at=%s elapsed_ms=%.2f",
|
||
step,
|
||
context.get("stat_date") or context.get("stat_dates") or "",
|
||
context.get("job_key") or "",
|
||
context.get("settings_hash") or "",
|
||
datetime.now().isoformat(),
|
||
(perf_counter() - start) * 1000,
|
||
)
|
||
return result
|
||
|
||
|
||
def _recent_leaderboard_stat_dates(days: int) -> list[str]:
|
||
from datetime import datetime, timedelta
|
||
|
||
from deerflow.persistence.types import BEIJING_TZ
|
||
|
||
today = datetime.now(BEIJING_TZ).date()
|
||
start = today - timedelta(days=max(1, days) - 1)
|
||
return [(start + timedelta(days=i)).isoformat() for i in range((today - start).days + 1)]
|
||
|
||
|
||
async def _leaderboard_snapshot_due_dates(
|
||
session_factory,
|
||
*,
|
||
settings_hash: str,
|
||
limit: int,
|
||
running_ttl_seconds: int,
|
||
) -> list[str]:
|
||
from datetime import UTC, datetime, timedelta
|
||
|
||
from sqlalchemy import and_, or_, select
|
||
|
||
from deerflow.persistence.admin_stats.model import AdminLeaderboardDailyStatRow
|
||
|
||
cutoff = datetime.now(UTC) - timedelta(seconds=running_ttl_seconds)
|
||
async with session_factory() as session:
|
||
stmt = (
|
||
select(AdminLeaderboardDailyStatRow.stat_date)
|
||
.where(
|
||
AdminLeaderboardDailyStatRow.settings_hash == settings_hash,
|
||
or_(
|
||
AdminLeaderboardDailyStatRow.status.in_(["queued", "error"]),
|
||
and_(
|
||
AdminLeaderboardDailyStatRow.status == "running",
|
||
AdminLeaderboardDailyStatRow.updated_at < cutoff,
|
||
),
|
||
),
|
||
)
|
||
.order_by(AdminLeaderboardDailyStatRow.stat_date.desc())
|
||
.limit(limit)
|
||
)
|
||
rows = (
|
||
await _execute_leaderboard_scheduler_query(
|
||
session,
|
||
"scheduler.due_dates",
|
||
stmt,
|
||
settings_hash=settings_hash,
|
||
)
|
||
).all()
|
||
return [str(row[0]) for row in rows]
|
||
|
||
|
||
async def _claim_leaderboard_snapshot(
|
||
session_factory,
|
||
*,
|
||
stat_date: str,
|
||
settings_hash: str,
|
||
running_ttl_seconds: int,
|
||
) -> bool:
|
||
from datetime import UTC, datetime
|
||
|
||
from app.gateway.routers.admin_users import _snapshot_is_fresh
|
||
from deerflow.persistence.admin_stats.model import AdminLeaderboardDailyStatRow
|
||
from deerflow.persistence.types import BEIJING_TZ
|
||
|
||
async with session_factory() as session:
|
||
row = await _time_leaderboard_scheduler_db_operation(
|
||
"scheduler.claim.get",
|
||
lambda: session.get(
|
||
AdminLeaderboardDailyStatRow,
|
||
{"stat_date": stat_date, "settings_hash": settings_hash},
|
||
),
|
||
sql="SELECT admin_leaderboard_daily_stats by primary key",
|
||
stat_date=stat_date,
|
||
settings_hash=settings_hash,
|
||
)
|
||
if row is None or _snapshot_is_fresh(row, stat_date):
|
||
return False
|
||
if row.status == "running":
|
||
updated_at = row.updated_at
|
||
if isinstance(updated_at, datetime):
|
||
if updated_at.tzinfo is None:
|
||
updated_at = updated_at.replace(tzinfo=BEIJING_TZ)
|
||
age = (datetime.now(BEIJING_TZ) - updated_at.astimezone(BEIJING_TZ)).total_seconds()
|
||
if age < running_ttl_seconds:
|
||
return False
|
||
row.status = "running"
|
||
row.message = "?????????????"
|
||
row.error = None
|
||
row.updated_at = datetime.now(UTC)
|
||
await _time_leaderboard_scheduler_db_operation(
|
||
"scheduler.claim.commit",
|
||
lambda: session.commit(),
|
||
sql="COMMIT",
|
||
stat_date=stat_date,
|
||
settings_hash=settings_hash,
|
||
)
|
||
return True
|
||
|
||
|
||
async def _maybe_finalize_previous_leaderboard_day(session_factory, settings_hash: str) -> bool:
|
||
"""Once per day after ``_LEADERBOARD_FINALIZE_HOUR``, force-recompute yesterday.
|
||
|
||
A ``ready`` past-day snapshot is otherwise treated as fresh forever and never
|
||
rebuilt. But the previous day's snapshot is usually generated before midnight,
|
||
so it can miss runs/tool-calls that only land in the DB around the day
|
||
boundary. This re-queues yesterday (``force=True``) a single time; the rebuild
|
||
stamps ``generated_at`` to today, which stops it from re-firing for the rest of
|
||
the day.
|
||
"""
|
||
from datetime import timedelta
|
||
|
||
from app.gateway.routers.admin_users import (
|
||
_as_beijing,
|
||
_queue_daily_leaderboard_snapshots,
|
||
)
|
||
from deerflow.persistence.admin_stats.finalize import should_finalize_previous_day
|
||
from deerflow.persistence.admin_stats.model import AdminLeaderboardDailyStatRow
|
||
from deerflow.persistence.types import BEIJING_TZ
|
||
|
||
now = datetime.now(BEIJING_TZ)
|
||
yesterday = (now.date() - timedelta(days=1)).isoformat()
|
||
async with session_factory() as session:
|
||
row = await session.get(
|
||
AdminLeaderboardDailyStatRow,
|
||
{"stat_date": yesterday, "settings_hash": settings_hash},
|
||
)
|
||
generated_at = _as_beijing(row.generated_at) if (row is not None and row.generated_at is not None) else None
|
||
if not should_finalize_previous_day(
|
||
status=(row.status if row is not None else None),
|
||
generated_at=generated_at,
|
||
now=now,
|
||
finalize_hour=_LEADERBOARD_FINALIZE_HOUR,
|
||
):
|
||
return False
|
||
queued = await _queue_daily_leaderboard_snapshots(
|
||
session_factory,
|
||
[yesterday],
|
||
settings_hash,
|
||
force=True,
|
||
message="???????????????????????????",
|
||
)
|
||
if queued:
|
||
logger.info(
|
||
"Leaderboard snapshot finalize: queued previous day %s for full recompute (hash=%s)",
|
||
yesterday,
|
||
settings_hash[:10],
|
||
)
|
||
return True
|
||
return False
|
||
|
||
|
||
async def _run_leaderboard_snapshot_tick() -> tuple[int, int]:
|
||
from app.gateway.routers.admin_users import (
|
||
_daily_job_key,
|
||
_queue_daily_leaderboard_snapshots,
|
||
_refresh_daily_leaderboard_snapshot,
|
||
_resolve_leaderboard_snapshot_scope,
|
||
)
|
||
from deerflow.config.system_settings import load_system_settings
|
||
from deerflow.persistence.engine import get_session_factory
|
||
|
||
session_factory = get_session_factory()
|
||
if session_factory is None:
|
||
return 0, 0
|
||
|
||
settings = load_system_settings().leaderboard
|
||
excluded, include_scheduled, include_admins, include_failed, settings_hash = _resolve_leaderboard_snapshot_scope(settings)
|
||
stat_dates = _recent_leaderboard_stat_dates(_LEADERBOARD_SNAPSHOT_PREWARM_DAYS)
|
||
queued_dates = await _queue_daily_leaderboard_snapshots(
|
||
session_factory,
|
||
stat_dates,
|
||
settings_hash,
|
||
message="?????????????/?????",
|
||
)
|
||
|
||
# ??????? 1 ???????"???"????????????????????
|
||
try:
|
||
await _maybe_finalize_previous_leaderboard_day(session_factory, settings_hash)
|
||
except Exception:
|
||
logger.exception("Leaderboard snapshot finalize check failed (non-fatal)")
|
||
|
||
due_dates = await _leaderboard_snapshot_due_dates(
|
||
session_factory,
|
||
settings_hash=settings_hash,
|
||
limit=_LEADERBOARD_SNAPSHOT_BATCH_SIZE,
|
||
running_ttl_seconds=_LEADERBOARD_SNAPSHOT_RUNNING_TTL_SECONDS,
|
||
)
|
||
processed = 0
|
||
for stat_date in due_dates:
|
||
claimed = await _claim_leaderboard_snapshot(
|
||
session_factory,
|
||
stat_date=stat_date,
|
||
settings_hash=settings_hash,
|
||
running_ttl_seconds=_LEADERBOARD_SNAPSHOT_RUNNING_TTL_SECONDS,
|
||
)
|
||
if not claimed:
|
||
continue
|
||
logger.info("Leaderboard snapshot scheduler: building date=%s hash=%s", stat_date, settings_hash[:10])
|
||
await _refresh_daily_leaderboard_snapshot(
|
||
_daily_job_key(stat_date, settings_hash),
|
||
session_factory=session_factory,
|
||
settings=settings,
|
||
stat_date=stat_date,
|
||
settings_hash=settings_hash,
|
||
excluded=excluded,
|
||
include_scheduled=include_scheduled,
|
||
include_admins=include_admins,
|
||
include_failed=include_failed,
|
||
)
|
||
processed += 1
|
||
if _LEADERBOARD_SNAPSHOT_SLEEP_SECONDS > 0 and processed < len(due_dates):
|
||
await asyncio.sleep(_LEADERBOARD_SNAPSHOT_SLEEP_SECONDS)
|
||
return processed, len(queued_dates)
|
||
|
||
|
||
async def _leaderboard_snapshot_scheduler_loop() -> None:
|
||
"""Background task inside the existing gateway process for daily snapshots."""
|
||
if not _LEADERBOARD_SNAPSHOT_BACKGROUND_ENABLED:
|
||
logger.info("Leaderboard snapshot scheduler disabled by ADMIN_LEADERBOARD_BACKGROUND_ENABLED")
|
||
return
|
||
|
||
from deerflow.config.runtime_paths import runtime_home
|
||
|
||
lock_handle = _acquire_singleton_lock(runtime_home() / ".leaderboard_snapshot_scheduler.lock")
|
||
if lock_handle is None:
|
||
logger.info("Leaderboard snapshot scheduler: another gateway worker holds the lock; this worker skips")
|
||
return
|
||
|
||
logger.info(
|
||
"Leaderboard snapshot scheduler started: prewarm_days=%d batch_size=%d interval=%.1fs",
|
||
_LEADERBOARD_SNAPSHOT_PREWARM_DAYS,
|
||
_LEADERBOARD_SNAPSHOT_BATCH_SIZE,
|
||
_LEADERBOARD_SNAPSHOT_INTERVAL_SECONDS,
|
||
)
|
||
try:
|
||
if _LEADERBOARD_SNAPSHOT_STARTUP_DELAY_SECONDS > 0:
|
||
await asyncio.sleep(_LEADERBOARD_SNAPSHOT_STARTUP_DELAY_SECONDS)
|
||
while True:
|
||
try:
|
||
processed, queued = await _run_leaderboard_snapshot_tick()
|
||
if processed or queued:
|
||
logger.info(
|
||
"Leaderboard snapshot scheduler tick: queued=%d processed=%d",
|
||
queued,
|
||
processed,
|
||
)
|
||
except Exception:
|
||
logger.exception("Leaderboard snapshot scheduler tick failed (non-fatal)")
|
||
await asyncio.sleep(_LEADERBOARD_SNAPSHOT_INTERVAL_SECONDS)
|
||
finally:
|
||
_release_singleton_lock(lock_handle)
|
||
|
||
|
||
async def _concurrency_sampler_loop(app: FastAPI) -> None:
|
||
"""Periodically record the in-flight run count for the admin pressure chart.
|
||
|
||
Runs in the gateway process, singleton-locked so only one worker samples.
|
||
Each tick records ``len(run_manager.list_active())`` and, roughly hourly,
|
||
purges samples older than the retention window. Best-effort throughout — a
|
||
failed tick is logged and the loop continues.
|
||
"""
|
||
if not _CONCURRENCY_SAMPLE_ENABLED:
|
||
logger.info("Concurrency sampler disabled by CONCURRENCY_SAMPLE_ENABLED")
|
||
return
|
||
|
||
store = getattr(app.state, "concurrency_sample_store", None)
|
||
run_mgr = getattr(app.state, "run_manager", None)
|
||
concurrency_gate = getattr(app.state, "concurrency_gate", None)
|
||
if store is None or run_mgr is None:
|
||
logger.info("Concurrency sampler not started (no store / run manager — memory backend?)")
|
||
return
|
||
|
||
from deerflow.config.runtime_paths import runtime_home
|
||
|
||
lock_handle = _acquire_singleton_lock(runtime_home() / ".concurrency_sampler.lock")
|
||
if lock_handle is None:
|
||
logger.info("Concurrency sampler: another gateway worker holds the lock; this worker skips")
|
||
return
|
||
|
||
interval = max(1.0, _CONCURRENCY_SAMPLE_INTERVAL_SECONDS)
|
||
purge_every_ticks = max(1, int(3600 / interval)) # ~once per hour
|
||
logger.info("Concurrency sampler started: interval=%.1fs retention=%.1fh", interval, _CONCURRENCY_SAMPLE_RETENTION_HOURS)
|
||
ticks = 0
|
||
try:
|
||
while True:
|
||
try:
|
||
# Runtime toggle (admin can turn sampling off without a restart).
|
||
# mtime-cached read, so this is cheap every tick.
|
||
from deerflow.config.system_settings import get_concurrency_monitor_settings
|
||
|
||
if get_concurrency_monitor_settings().enabled:
|
||
active_count = None
|
||
if concurrency_gate is not None:
|
||
active_count = await concurrency_gate.count_active()
|
||
if active_count is None:
|
||
active = await run_mgr.list_active()
|
||
active_count = len(active)
|
||
await store.record(active_count)
|
||
ticks += 1
|
||
if ticks % purge_every_ticks == 0:
|
||
from datetime import UTC, datetime, timedelta
|
||
|
||
cutoff = datetime.now(UTC) - timedelta(hours=_CONCURRENCY_SAMPLE_RETENTION_HOURS)
|
||
await store.purge_older_than(cutoff)
|
||
except Exception:
|
||
logger.exception("Concurrency sampler tick failed (non-fatal)")
|
||
await asyncio.sleep(interval)
|
||
finally:
|
||
_release_singleton_lock(lock_handle)
|
||
|
||
|
||
def _warn_if_multi_worker() -> None:
|
||
"""?? uvicorn / gunicorn ? worker ???? ERROR ???
|
||
|
||
? agent / canvas / AI ??????? Gateway ????? ``StreamBridge``
|
||
?? run_id ? SSE ?????? map?? worker ????? worker A ??
|
||
``/runs/stream``???? worker B ? join/resume??? in-process map ??
|
||
???? 404 / ???
|
||
|
||
????? startup ? ERROR ????**???? fail-fast** ?? ??
|
||
???????????????? sticky session??????????
|
||
"""
|
||
# uvicorn / gunicorn ??????? worker ?????????????
|
||
# ?????? sibling ??????????? WEB_CONCURRENCY?gunicorn
|
||
# / uvicorn ????
|
||
web_concurrency = os.environ.get("WEB_CONCURRENCY", "").strip()
|
||
if web_concurrency and web_concurrency.isdigit() and int(web_concurrency) > 1:
|
||
logger.error(
|
||
"? ??? WEB_CONCURRENCY=%s?? worker ????AI ?? / ? agent "
|
||
"??????????????? worker ????????? / "
|
||
"SSE ?? / resume 400??????????? WEB_CONCURRENCY=1 "
|
||
"????????? sticky session??? AIWritingPage.tsx ????",
|
||
web_concurrency,
|
||
)
|
||
return
|
||
|
||
# ???? /proc ??? uvicorn ?????? Linux ??????
|
||
if not sys.platform.startswith("linux"):
|
||
return
|
||
try:
|
||
from pathlib import Path
|
||
|
||
peers = 0
|
||
for entry in Path("/proc").iterdir():
|
||
if not entry.name.isdigit():
|
||
continue
|
||
try:
|
||
cmdline = (entry / "cmdline").read_bytes().replace(b"\x00", b" ").decode("utf-8", "ignore")
|
||
except OSError:
|
||
continue
|
||
if "uvicorn" in cmdline and "app.gateway.app" in cmdline:
|
||
peers += 1
|
||
if peers > 1:
|
||
logger.error(
|
||
"? ??? %d ??? uvicorn ???? worker ????? AI ?? "
|
||
"????? / SSE ??????????? worker ??????? "
|
||
"?? sticky session?",
|
||
peers,
|
||
)
|
||
except Exception:
|
||
# ????????? ?? ??????
|
||
pass
|
||
|
||
|
||
@asynccontextmanager
|
||
async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
|
||
"""Application lifespan handler."""
|
||
|
||
# ?????? worker ?? ?? ?? AI ?? / ? agent ???????
|
||
# ??????? worker ???????? ERROR ?????
|
||
_warn_if_multi_worker()
|
||
|
||
# Load config and check necessary environment variables at startup
|
||
try:
|
||
app.state.config = get_app_config()
|
||
apply_logging_level(app.state.config.log_level)
|
||
logger.info("Configuration loaded successfully")
|
||
except Exception as e:
|
||
error_msg = f"Failed to load configuration during gateway startup: {e}"
|
||
logger.exception(error_msg)
|
||
raise RuntimeError(error_msg) from e
|
||
config = get_gateway_config()
|
||
logger.info(f"Starting API Gateway on {config.host}:{config.port}")
|
||
|
||
# Initialize LangGraph runtime components (StreamBridge, RunManager, checkpointer, store)
|
||
async with langgraph_runtime(app):
|
||
logger.info("LangGraph runtime initialised")
|
||
|
||
# One-time repair switch for a failed, unused Workflow Studio schema.
|
||
# The helper is intentionally MySQL-only and drops an exact allow-list
|
||
# in the configured current database; it never expands ``workflow_%``.
|
||
# Keep this opt-in so established workflow data is never cleared during
|
||
# an ordinary restart.
|
||
if os.getenv("DEER_FLOW_REBUILD_WORKFLOW_SCHEMA", "").lower() in {"1", "true", "yes"}:
|
||
from deerflow.persistence.engine import rebuild_workflow_schema_for_mysql
|
||
|
||
logger.warning(
|
||
"DEER_FLOW_REBUILD_WORKFLOW_SCHEMA is enabled; rebuilding the unused workflow schema"
|
||
)
|
||
await rebuild_workflow_schema_for_mysql()
|
||
|
||
# Run Alembic migrations before anything else touches the schema.
|
||
await _run_db_migrations()
|
||
|
||
# Local Wiki mirror/index scans are opt-in and run outside request
|
||
# handling. Per-library database leases keep multi-worker startup safe.
|
||
app.state.llmwiki_index_scheduler_task = None
|
||
local_wiki_index = app.state.config.llmwiki.local_wiki_index
|
||
if (
|
||
local_wiki_index.enabled
|
||
and local_wiki_index.auto_sync
|
||
and getattr(app.state, "llmwiki_sync_service", None) is not None
|
||
):
|
||
from app.gateway.llmwiki_index_scheduler import local_wiki_index_scheduler_loop
|
||
|
||
app.state.llmwiki_index_scheduler_task = asyncio.create_task(
|
||
local_wiki_index_scheduler_loop(app)
|
||
)
|
||
|
||
# ?? recall_budget cache?? DB ????,? get_effective_memory_config ?????
|
||
try:
|
||
from deerflow.config.memory_config import load_recall_budget_cache
|
||
|
||
_mc_store = getattr(app.state, "memory_config_store", None)
|
||
if _mc_store is not None:
|
||
_recall_budgets = await _mc_store.load_all()
|
||
load_recall_budget_cache(_recall_budgets)
|
||
logger.info("Memory config cache warmed: %d user override(s)", len(_recall_budgets))
|
||
except Exception:
|
||
logger.exception("Memory config cache warm-up failed (non-fatal)")
|
||
|
||
# Ensure admin user exists (auto-create on first boot)
|
||
# Must run AFTER langgraph_runtime so app.state.store is available for thread migration
|
||
await _ensure_admin_user(app)
|
||
|
||
# Migrate "ownerless" on-disk skills into the new SkillStore table ?
|
||
# makes the multitenancy switch transparent for existing deployments.
|
||
try:
|
||
await _sync_legacy_skills(app)
|
||
except Exception:
|
||
logger.exception("Skill ownership migration failed (non-fatal)")
|
||
|
||
# ????? agent(intent / recommender / coordinator)???????
|
||
# ??? _sync_legacy_agents **??**??????????????????
|
||
# ??? upsert ? agents ?,?????? /api/agents ?????,
|
||
# ????????????????.deer-flow/ ? gitignore,??? clone
|
||
# ????????,seeder ????? _roundtable_seed_assets/ ??????
|
||
try:
|
||
from app.gateway.routers._roundtable_seed import (
|
||
ensure_roundtable_functional_agents,
|
||
)
|
||
|
||
ensure_roundtable_functional_agents()
|
||
except Exception:
|
||
logger.exception(
|
||
"Roundtable functional agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# AI ????? agent(researcher / outliner / writer / editor)???????
|
||
# ???:??? _sync_legacy_agents ??,?????????????
|
||
# upsert ? agents ?,admin ?????????????????? agent?
|
||
try:
|
||
from app.gateway.routers._ai_writing_seed import (
|
||
ensure_ai_writing_functional_agents,
|
||
)
|
||
|
||
ensure_ai_writing_functional_agents()
|
||
except Exception:
|
||
logger.exception(
|
||
"AI-writing functional agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# 内置「模板数据转换器」agent(JSON → HTML 模板填充)。同样必须在
|
||
# _sync_legacy_agents **之前**补建目录,后者把它 upsert 进 agents 表,
|
||
# 第一次启动 /api/agents 与定时任务下拉框里就有它。
|
||
try:
|
||
from app.gateway.routers._template_builder_seed import (
|
||
ensure_template_builder_agent,
|
||
)
|
||
|
||
ensure_template_builder_agent()
|
||
except Exception:
|
||
logger.exception(
|
||
"Template-builder agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# 实验性「页面研报设计师」:在 _sync_legacy_agents 前补齐 agent 目录与
|
||
# 公共页面生成技能,首次启动即可在智能体管理页和普通对话中复用。
|
||
try:
|
||
from app.gateway.routers._page_report_designer_seed import (
|
||
ensure_page_report_designer_agent,
|
||
)
|
||
|
||
ensure_page_report_designer_agent()
|
||
except Exception:
|
||
logger.exception(
|
||
"Page-report designer agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# 深度研究「收集员」内置 agent(对话式信息收集)。同样必须在
|
||
# _sync_legacy_agents 之前补建目录,由后者 upsert 进 agents 表,
|
||
# 深度研究页的收集对话与智能体管理页即可直接使用。
|
||
try:
|
||
from app.gateway.routers._deep_research_seed import (
|
||
ensure_deep_research_collector_agent,
|
||
)
|
||
|
||
ensure_deep_research_collector_agent()
|
||
except Exception:
|
||
logger.exception(
|
||
"Deep-research collector agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# Workflow Studio owns a separate planning controller. It must be
|
||
# seeded before _sync_legacy_agents, so the planner is available to
|
||
# every user's candidate-workflow request on a fresh deployment.
|
||
try:
|
||
from app.gateway.routers._workflow_planner_seed import (
|
||
ensure_workflow_planner_agent,
|
||
)
|
||
|
||
ensure_workflow_planner_agent()
|
||
except Exception:
|
||
logger.exception(
|
||
"Workflow planner agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# 内置「强制检索输出助手」:必须在 _sync_legacy_agents 之前补建,
|
||
# 使首次启动就能在智能体管理/广场中使用。每轮检索门禁由 harness
|
||
# 的 ForcedResearchMiddleware 执行,而不是仅依赖 SOUL 提示词。
|
||
try:
|
||
from app.gateway.routers._forced_research_seed import (
|
||
ensure_forced_research_agent,
|
||
)
|
||
|
||
ensure_forced_research_agent()
|
||
except Exception:
|
||
logger.exception(
|
||
"Forced-research responder agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# 内置「任务研判报告助手」(agentfx 任务深链固定使用):同样必须在
|
||
# _sync_legacy_agents 之前补建目录 + 补建配套 task-report-import 技能
|
||
# (skills/ 被 .gitignore 排除),由后者 upsert 进 agents 表。
|
||
try:
|
||
from app.gateway.routers._agentfx_seed import ensure_agentfx_agent
|
||
|
||
ensure_agentfx_agent()
|
||
except Exception:
|
||
logger.exception(
|
||
"Agentfx analyst agent seeding failed at startup (non-fatal)"
|
||
)
|
||
|
||
# ?????????????????3Q/6BF/7BF/8BF??agent/chain ?????
|
||
try:
|
||
business_mapping_store = getattr(app.state, "business_mapping_store", None)
|
||
if business_mapping_store is not None:
|
||
inserted = await business_mapping_store.ensure_seed_defaults()
|
||
if inserted:
|
||
logger.info("Seeded %d business mapping rows", inserted)
|
||
except Exception:
|
||
logger.exception("Business mapping seeding failed at startup (non-fatal)")
|
||
|
||
# ???? .deer-flow/agents/{id}/ ???????? agent ?? agents ??
|
||
# ??????????????? sync(????? + N ? agent ??? DB
|
||
# SELECT/UPDATE),??????? + ?????? TTL ??,????
|
||
# /api/agents ?????
|
||
try:
|
||
from app.gateway.routers.agents import (
|
||
_mark_legacy_sync_fresh,
|
||
_sync_legacy_agents,
|
||
)
|
||
from app.gateway.services import _mark_legacy_sync_fresh as _services_mark_sync_fresh
|
||
|
||
agent_store = getattr(app.state, "agent_store", None)
|
||
if agent_store is not None:
|
||
await _sync_legacy_agents(agent_store)
|
||
_mark_legacy_sync_fresh()
|
||
# ``services.py`` keeps an independent TTL gate (used by
|
||
# ``start_run`` ? ``_enforce_agent_access``). Without this
|
||
# ack, the first ``/runs/stream`` after startup would
|
||
# re-run the heavy sync even though we just did it.
|
||
_services_mark_sync_fresh()
|
||
logger.info("Legacy folder-backed agents synced at startup")
|
||
# 内置功能型 agent(roundtable-* / ai-writing-*)按用途默认打标签,
|
||
# 供管理员「内置智能体」筛选按用途归类。幂等、只追加、失败不阻塞启动。
|
||
try:
|
||
from app.gateway.routers._builtin_agent_tag_seed import (
|
||
ensure_builtin_agent_purpose_tags,
|
||
)
|
||
|
||
tag_store = getattr(app.state, "tag_store", None)
|
||
if tag_store is not None:
|
||
await ensure_builtin_agent_purpose_tags(agent_store, tag_store)
|
||
except Exception:
|
||
logger.exception("Built-in agent purpose-tag seeding failed at startup (non-fatal)")
|
||
except Exception:
|
||
logger.exception("Legacy agent sync failed at startup (non-fatal)")
|
||
|
||
# ?? skill curator ???:??????,?????????????
|
||
# (per-user interval_hours ??????????????)?
|
||
try:
|
||
app.state.curator_scheduler_task = asyncio.create_task(_curator_scheduler_loop())
|
||
except Exception:
|
||
logger.exception("Skill curator scheduler failed to start (non-fatal)")
|
||
|
||
# ?? AI ????????:? config.yaml ? ai_writing.cleanup ???,
|
||
# ?????? 3 ????,? retention_days(?? 7 ?)?????,??
|
||
# ?? MySQL ??? + SQLite checkpointer ?? thread ???? worker
|
||
# ?????? leader,?????? worker ??????
|
||
try:
|
||
from app.gateway.ai_writing_cleanup import cleanup_scheduler_loop
|
||
|
||
app.state.ai_writing_cleanup_task = asyncio.create_task(cleanup_scheduler_loop(app))
|
||
except Exception:
|
||
logger.exception("AI writing cleanup scheduler failed to start (non-fatal)")
|
||
|
||
try:
|
||
app.state.leaderboard_snapshot_scheduler_task = asyncio.create_task(_leaderboard_snapshot_scheduler_loop())
|
||
except Exception:
|
||
logger.exception("Leaderboard snapshot scheduler failed to start (non-fatal)")
|
||
|
||
# Concurrency (system-pressure) sampler — lightweight, singleton-locked,
|
||
# runtime-toggleable via system_settings.concurrency_monitor.enabled.
|
||
try:
|
||
app.state.concurrency_sampler_task = asyncio.create_task(_concurrency_sampler_loop(app))
|
||
except Exception:
|
||
logger.exception("Concurrency sampler failed to start (non-fatal)")
|
||
|
||
# Start IM channel service if any channels are configured
|
||
try:
|
||
from app.channels.service import start_channel_service
|
||
|
||
channel_service = await start_channel_service(app.state.config)
|
||
logger.info("Channel service started: %s", channel_service.get_status())
|
||
except Exception:
|
||
logger.exception("No IM channels configured or channel service failed to start")
|
||
|
||
scheduled_tasks_enabled = os.environ.get("SCHEDULED_TASKS_ENABLED", "true").strip().lower()
|
||
if scheduled_tasks_enabled in {"0", "false", "no", "off"}:
|
||
logger.info("Scheduled task service disabled by SCHEDULED_TASKS_ENABLED=%s", scheduled_tasks_enabled)
|
||
else:
|
||
try:
|
||
scheduled_task_store = getattr(app.state, "scheduled_task_store", None)
|
||
if scheduled_task_store is not None:
|
||
scheduler_service = ScheduledTaskService(app, scheduled_task_store)
|
||
app.state.scheduled_task_service = scheduler_service
|
||
set_scheduled_task_service(scheduler_service)
|
||
await scheduler_service.start()
|
||
logger.info("Scheduled task service started")
|
||
except Exception:
|
||
logger.exception("Scheduled task service failed to start")
|
||
|
||
yield
|
||
|
||
llmwiki_index_scheduler_task = getattr(app.state, "llmwiki_index_scheduler_task", None)
|
||
if llmwiki_index_scheduler_task is not None:
|
||
llmwiki_index_scheduler_task.cancel()
|
||
try:
|
||
await llmwiki_index_scheduler_task
|
||
except asyncio.CancelledError:
|
||
pass
|
||
except Exception:
|
||
logger.exception("Local Wiki index scheduler shutdown error")
|
||
llmwiki_index_jobs = list(getattr(app.state, "llmwiki_index_jobs", set()))
|
||
for job in llmwiki_index_jobs:
|
||
job.cancel()
|
||
if llmwiki_index_jobs:
|
||
await asyncio.gather(*llmwiki_index_jobs, return_exceptions=True)
|
||
|
||
# ?? curator ???
|
||
curator_task = getattr(app.state, "curator_scheduler_task", None)
|
||
if curator_task is not None:
|
||
curator_task.cancel()
|
||
try:
|
||
await curator_task
|
||
except asyncio.CancelledError:
|
||
pass
|
||
except Exception:
|
||
logger.exception("Curator scheduler shutdown error")
|
||
|
||
# ?? AI ???????
|
||
cleanup_task = getattr(app.state, "ai_writing_cleanup_task", None)
|
||
if cleanup_task is not None:
|
||
cleanup_task.cancel()
|
||
try:
|
||
await cleanup_task
|
||
except asyncio.CancelledError:
|
||
pass
|
||
except Exception:
|
||
logger.exception("AI writing cleanup scheduler shutdown error")
|
||
|
||
leaderboard_snapshot_task = getattr(app.state, "leaderboard_snapshot_scheduler_task", None)
|
||
if leaderboard_snapshot_task is not None:
|
||
leaderboard_snapshot_task.cancel()
|
||
try:
|
||
await leaderboard_snapshot_task
|
||
except asyncio.CancelledError:
|
||
pass
|
||
except Exception:
|
||
logger.exception("Leaderboard snapshot scheduler shutdown error")
|
||
|
||
concurrency_sampler_task = getattr(app.state, "concurrency_sampler_task", None)
|
||
if concurrency_sampler_task is not None:
|
||
concurrency_sampler_task.cancel()
|
||
try:
|
||
await concurrency_sampler_task
|
||
except asyncio.CancelledError:
|
||
pass
|
||
except Exception:
|
||
logger.exception("Concurrency sampler shutdown error")
|
||
|
||
# ?? 7 ????AI ?? graph ???? checkpointer / in-flight task?
|
||
# thread ??? LangGraph Server ? checkpointer ???lifespan ???
|
||
# ``langgraph_runtime()`` ? contextmanager ???????????????
|
||
|
||
try:
|
||
scheduler_service = getattr(app.state, "scheduled_task_service", None)
|
||
if scheduler_service is not None:
|
||
await scheduler_service.stop()
|
||
except Exception:
|
||
logger.exception("Failed to stop scheduled task service")
|
||
finally:
|
||
set_scheduled_task_service(None)
|
||
|
||
# Stop channel service on shutdown (bounded to prevent worker hang)
|
||
try:
|
||
from app.channels.service import stop_channel_service
|
||
|
||
await asyncio.wait_for(
|
||
stop_channel_service(),
|
||
timeout=_SHUTDOWN_HOOK_TIMEOUT_SECONDS,
|
||
)
|
||
except TimeoutError:
|
||
logger.warning(
|
||
"Channel service shutdown exceeded %.1fs; proceeding with worker exit.",
|
||
_SHUTDOWN_HOOK_TIMEOUT_SECONDS,
|
||
)
|
||
except Exception:
|
||
logger.exception("Failed to stop channel service")
|
||
|
||
logger.info("Shutting down API Gateway")
|
||
|
||
|
||
def create_app() -> FastAPI:
|
||
"""Create and configure the FastAPI application.
|
||
|
||
Returns:
|
||
Configured FastAPI application instance.
|
||
"""
|
||
config = get_gateway_config()
|
||
docs_kwargs = {"docs_url": "/docs", "redoc_url": "/redoc", "openapi_url": "/openapi.json"} if config.enable_docs else {"docs_url": None, "redoc_url": None, "openapi_url": None}
|
||
|
||
app = FastAPI(
|
||
title="DeerFlow API Gateway",
|
||
description="""
|
||
## DeerFlow API Gateway
|
||
|
||
API Gateway for DeerFlow - A LangGraph-based AI agent backend with sandbox execution capabilities.
|
||
|
||
### Features
|
||
|
||
- **Models Management**: Query and retrieve available AI models
|
||
- **MCP Configuration**: Manage Model Context Protocol (MCP) server configurations
|
||
- **Memory Management**: Access and manage global memory data for personalized conversations
|
||
- **Skills Management**: Query and manage skills and their enabled status
|
||
- **Artifacts**: Access thread artifacts and generated files
|
||
- **Health Monitoring**: System health check endpoints
|
||
|
||
### Architecture
|
||
|
||
LangGraph requests are handled by nginx reverse proxy.
|
||
This gateway provides custom endpoints for models, MCP configuration, skills, and artifacts.
|
||
""",
|
||
version="0.1.0",
|
||
lifespan=lifespan,
|
||
**docs_kwargs,
|
||
openapi_tags=[
|
||
{
|
||
"name": "models",
|
||
"description": "Operations for querying available AI models and their configurations",
|
||
},
|
||
{
|
||
"name": "mcp",
|
||
"description": "Manage Model Context Protocol (MCP) server configurations",
|
||
},
|
||
{
|
||
"name": "memory",
|
||
"description": "Access and manage global memory data for personalized conversations",
|
||
},
|
||
{
|
||
"name": "skills",
|
||
"description": "Manage skills and their configurations",
|
||
},
|
||
{
|
||
"name": "artifacts",
|
||
"description": "Access and download thread artifacts and generated files",
|
||
},
|
||
{
|
||
"name": "uploads",
|
||
"description": "Upload and manage user files for threads",
|
||
},
|
||
{
|
||
"name": "threads",
|
||
"description": "Manage DeerFlow thread-local filesystem data",
|
||
},
|
||
{
|
||
"name": "agents",
|
||
"description": "Create and manage custom agents with per-agent config and prompts",
|
||
},
|
||
{
|
||
"name": "suggestions",
|
||
"description": "Generate follow-up question suggestions for conversations",
|
||
},
|
||
{
|
||
"name": "channels",
|
||
"description": "Manage IM channel integrations (Feishu, Slack, Telegram)",
|
||
},
|
||
{
|
||
"name": "assistants-compat",
|
||
"description": "LangGraph Platform-compatible assistants API (stub)",
|
||
},
|
||
{
|
||
"name": "runs",
|
||
"description": "LangGraph Platform-compatible runs lifecycle (create, stream, cancel)",
|
||
},
|
||
{
|
||
"name": "scheduled-tasks",
|
||
"description": "Create, manage, and execute user scheduled agent tasks",
|
||
},
|
||
{
|
||
"name": "health",
|
||
"description": "Health check and system status endpoints",
|
||
},
|
||
],
|
||
)
|
||
|
||
# Auth: reject unauthenticated requests to non-public paths (fail-closed safety net)
|
||
app.add_middleware(AuthMiddleware)
|
||
|
||
# CSRF: Double Submit Cookie pattern for state-changing requests
|
||
app.add_middleware(CSRFMiddleware)
|
||
|
||
# CORS
|
||
if is_auth_disabled():
|
||
logger.warning(
|
||
"DEER_FLOW_AUTH_DISABLED is set ? API authentication is bypassed. "
|
||
"Do not use this in production."
|
||
)
|
||
|
||
def _cors_allow_any_origin() -> bool:
|
||
"""Mirror any browser Origin (works with credentials); cannot use literal '*'."""
|
||
v = os.environ.get("GATEWAY_CORS_ALLOW_ALL", "").strip().lower()
|
||
return v in ("1", "true", "yes", "on") or is_auth_disabled()
|
||
|
||
# Keep standard local frontend origins available even when the startup
|
||
# script supplies an explicit production/development allow-list. Vite can
|
||
# choose 5173 or move to 5174 when that port is already occupied; excluding
|
||
# the former makes every browser preflight fail after a backend restart.
|
||
local_vite_origins = [
|
||
"http://localhost:5173",
|
||
"http://127.0.0.1:5173",
|
||
"http://localhost:5174",
|
||
"http://127.0.0.1:5174",
|
||
"http://localhost:3000",
|
||
"http://127.0.0.1:3000",
|
||
]
|
||
|
||
# The legacy TaskCOP Vue application uses its own Vite development port.
|
||
taskcop_vite_origins = [
|
||
"http://localhost:8080",
|
||
"http://127.0.0.1:8080",
|
||
]
|
||
local_dev_origins = [*local_vite_origins, *taskcop_vite_origins]
|
||
|
||
cors_origins_env = os.environ.get("GATEWAY_CORS_ORIGINS", "").strip()
|
||
cors_origins: list[str] = [o.strip() for o in cors_origins_env.split(",") if o.strip()] if cors_origins_env else []
|
||
if cors_origins:
|
||
cors_origins = list(dict.fromkeys([*cors_origins, *local_dev_origins]))
|
||
|
||
# Headers the LangGraph SDK / browser clients need to read from cross-origin
|
||
# responses. Without exposing them, the SDK's `response.headers.get(...)` returns
|
||
# null on cross-origin streaming responses, breaking run-id extraction and reconnect.
|
||
cors_expose_headers = [
|
||
"Content-Disposition",
|
||
"Content-Location",
|
||
"Location",
|
||
"X-Run-Id",
|
||
"X-Thread-Id",
|
||
]
|
||
|
||
# Public pages (e.g. http://47.x.x.x:7010) calling this Gateway on localhost
|
||
# trigger Chrome/Edge Private Network Access preflights. Without
|
||
# allow_private_network=True the browser reports a generic CORS failure even
|
||
# when Access-Control-Allow-Origin is correct.
|
||
cors_kwargs = {
|
||
"allow_credentials": True,
|
||
"allow_methods": ["*"],
|
||
"allow_headers": ["*"],
|
||
"expose_headers": cors_expose_headers,
|
||
"allow_private_network": True,
|
||
}
|
||
|
||
if _cors_allow_any_origin():
|
||
logger.warning(
|
||
"Permissive CORS enabled ? any Origin is echoed back (GATEWAY_CORS_ALLOW_ALL or "
|
||
"DEER_FLOW_AUTH_DISABLED). Do not expose publicly."
|
||
)
|
||
app.add_middleware(
|
||
CORSMiddleware,
|
||
allow_origins=[],
|
||
# fullmatch() on the Origin header; browsers send scheme://host[:port]
|
||
allow_origin_regex=r".+",
|
||
**cors_kwargs,
|
||
)
|
||
elif cors_origins:
|
||
if "*" in cors_origins:
|
||
logger.error(
|
||
"GATEWAY_CORS_ORIGINS contains '*' with allow_credentials=True ? invalid. "
|
||
"Set GATEWAY_CORS_ALLOW_ALL=1 instead, or list explicit origins."
|
||
)
|
||
cors_origins = [o for o in cors_origins if o != "*"]
|
||
if cors_origins:
|
||
app.add_middleware(
|
||
CORSMiddleware,
|
||
allow_origins=cors_origins,
|
||
**cors_kwargs,
|
||
)
|
||
else:
|
||
# Safe-by-default fallback for local frontend development. Without any CORS
|
||
# middleware, browser preflight OPTIONS requests will hit routes directly and
|
||
# fail with 405, which blocks login and all JSON POST APIs.
|
||
dev_origins = local_dev_origins
|
||
logger.warning(
|
||
"No CORS configuration found; enabling local dev origins only: %s. "
|
||
"Set GATEWAY_CORS_ORIGINS or GATEWAY_CORS_ALLOW_ALL for deployment.",
|
||
", ".join(dev_origins),
|
||
)
|
||
app.add_middleware(
|
||
CORSMiddleware,
|
||
allow_origins=dev_origins,
|
||
**cors_kwargs,
|
||
)
|
||
|
||
# This one anonymous API is intentionally callable by browser code from
|
||
# every origin. Keep the override route-scoped so authenticated APIs retain
|
||
# the deployment's normal CORS allow-list.
|
||
app.add_middleware(PublicKnowledgeCORSMiddleware)
|
||
|
||
# Include routers
|
||
# Models API is mounted at /api/models
|
||
app.include_router(models.router)
|
||
# Administrator-managed config.yaml model settings and connection tests.
|
||
app.include_router(model_management.router)
|
||
|
||
# MCP API is mounted at /api/mcp
|
||
app.include_router(mcp.router)
|
||
|
||
# DeerFlow-native LLMWiki management/search BFF (conditionally backed by WeKnora).
|
||
app.include_router(llmwiki.router)
|
||
app.include_router(llmwiki_index.router)
|
||
app.include_router(external_llmwiki.router)
|
||
|
||
# Memory API is mounted at /api/memory
|
||
app.include_router(memory.router)
|
||
|
||
# Skills API is mounted at /api/skills
|
||
app.include_router(skills.router)
|
||
# Admin skill distillation matrix + global assistant knowledge base.
|
||
app.include_router(skill_knowledge.router)
|
||
app.include_router(assistant_knowledge.router)
|
||
app.include_router(assistant_knowledge.export_router)
|
||
from app.gateway.routers import wiki_packages
|
||
|
||
app.include_router(wiki_packages.router)
|
||
|
||
# Artifacts API is mounted at /api/threads/{thread_id}/artifacts
|
||
app.include_router(artifacts.router)
|
||
|
||
# Uploads API is mounted at /api/threads/{thread_id}/uploads
|
||
app.include_router(uploads.router)
|
||
|
||
# Optional Chrome-extension page capture API is mounted at /api/browser-context
|
||
app.include_router(browser_context.router)
|
||
|
||
# Thread cleanup API is mounted at /api/threads/{thread_id}
|
||
app.include_router(threads.router)
|
||
|
||
# Agents API is mounted at /api/agents
|
||
app.include_router(agents.router)
|
||
|
||
# Recommended Questions API is mounted at /api/recommended-questions
|
||
app.include_router(recommended_questions.router)
|
||
|
||
app.include_router(user_prompts.router)
|
||
|
||
# Per-user UI preferences at /api/user-preferences
|
||
app.include_router(user_preferences.router)
|
||
|
||
# Roundtable-planning drafts + recommendation history at /api/roundtable-drafts
|
||
app.include_router(roundtable_drafts.router)
|
||
|
||
# Task-scoped roundtable drafts (taskId 深链聊天记录,独立存、不分权) at /api/roundtable-task-drafts
|
||
app.include_router(roundtable_task_drafts.router)
|
||
|
||
# Task-scoped artifact submission status; stores references only, never artifact content.
|
||
app.include_router(roundtable_artifact_submissions.router)
|
||
|
||
# Step3 流程图结构化抽取 + 校验 (SSE) at /api/flow-extract
|
||
app.include_router(flow_extract.router)
|
||
|
||
# Roundtable-planning business chains at /api/roundtable-chains
|
||
app.include_router(roundtable_chains.router)
|
||
|
||
# 岗位协同独立会话 / 节点状态(不使用原圆桌总控或后台作业)
|
||
# Canonical position-workspace persistence API.
|
||
app.include_router(position_roundtable.router, prefix="/api")
|
||
# Compatibility route for intranet gateways whose allowlist only exposes
|
||
# the established `/api/multi-agent/*` namespace.
|
||
app.include_router(
|
||
position_roundtable.router,
|
||
prefix="/api/multi-agent",
|
||
include_in_schema=False,
|
||
)
|
||
|
||
# Global business ? agent / chain mapping at /api/business-mapping
|
||
app.include_router(business_mapping.router)
|
||
|
||
# Position-roundtable work-view role configuration at /api/position-roles.
|
||
app.include_router(position_roles.router)
|
||
|
||
# Light-app registry (?????) at /api/light-apps
|
||
app.include_router(light_apps.router)
|
||
|
||
# Workflow Studio: definitions / runs / SSE / data sources / embed / Coze compat.
|
||
# Static prefixes (node-types, resources, planning, runs, data-sources, embed) must be
|
||
# registered before the /{workflow_id} catch-alls in workflows.router.
|
||
app.include_router(workflow_resources.router)
|
||
app.include_router(workflow_data_sources.router)
|
||
app.include_router(workflow_embed.router)
|
||
app.include_router(workflow_planning.router)
|
||
app.include_router(workflow_runs.router)
|
||
app.include_router(workflow_studio.router)
|
||
app.include_router(workflows.router)
|
||
app.include_router(workflows_coze_compat.router)
|
||
app.include_router(workflows_coze_compat.chrome_router)
|
||
|
||
# Sidebar-menu overrides (菜单管理) at /api/menu-overrides
|
||
app.include_router(menu_overrides.router)
|
||
|
||
# Compact/sidebar-summary menu overrides at /api/compact-menu-overrides
|
||
app.include_router(compact_menu_overrides.router)
|
||
|
||
# Sensitive-word map (敏感词管理 / 脱敏词) at /api/sensitive-words
|
||
app.include_router(sensitive_words.router)
|
||
|
||
# Sentiment-analysis virtual agent AG-UI SSE proxy at /api/sentiment-agent
|
||
app.include_router(sentiment_agent.router)
|
||
|
||
# Skill URL/IP batch replace (技能地址管理) at /api/skill-addresses
|
||
app.include_router(skill_addresses.router)
|
||
|
||
# 大屏绘制智能体独立页面会话(按 taskId 存取,不分权)at /api/dashboard-sessions
|
||
app.include_router(dashboard_sessions.router)
|
||
|
||
# Task-workspace configurable jump buttons (按钮管理) at /api/task-buttons
|
||
app.include_router(task_buttons.router)
|
||
|
||
# Admin exact-match question/answer short-circuit at /api/fixed-questions
|
||
app.include_router(fixed_questions.router)
|
||
|
||
# AgentScope report-collaboration workbench at /api/report-collaboration
|
||
app.include_router(report_collaboration.router)
|
||
|
||
# Deep-research report structures (工作流配置 / 报告结构) at /api/report-structures
|
||
app.include_router(report_structures.router)
|
||
|
||
# Temporary compatibility bridge for the TaskCOP situation-overview frontend.
|
||
app.include_router(taskcop_tasks.router)
|
||
|
||
# Task-scoped situation-report details used by dashboard drawing / task gates.
|
||
app.include_router(task_reports.router)
|
||
|
||
# Standalone parallel multi-agent panel: fan-out streaming + saved tasks
|
||
app.include_router(parallel_agents.router)
|
||
app.include_router(parallel_agent_tasks.router)
|
||
|
||
# Roundtable-planning background jobs at /api/roundtable-jobs
|
||
app.include_router(roundtable_jobs.router)
|
||
app.include_router(roundtable_diagnostics.router)
|
||
|
||
# Deep Research (深度研究) API at /api/deep-research
|
||
app.include_router(deep_research.router)
|
||
|
||
# The independent enterprise-research workbench is temporarily disabled.
|
||
# Do not mount /api/enterprise-research while its MySQL schema is inactive.
|
||
|
||
# Suggestions API is mounted at /api/threads/{thread_id}/suggestions
|
||
app.include_router(suggestions.router)
|
||
|
||
# Writing rewrite API is mounted at /api/writing
|
||
app.include_router(writing.router)
|
||
|
||
# AI Writing multi-agent API is mounted at /api/ai-writing
|
||
app.include_router(ai_writing.router)
|
||
|
||
# Scheduled task API is mounted at /api/scheduled-tasks
|
||
app.include_router(scheduled_tasks.router)
|
||
|
||
# Public, unauthenticated read-only scheduled task API
|
||
# (/api/public/scheduled-tasks/*). Lets external systems query tasks and
|
||
# surface results without a DeerFlow login.
|
||
app.include_router(public_scheduled_tasks.router)
|
||
|
||
# Public read-only user leaderboard publish page (/api/public/leaderboard)
|
||
app.include_router(public_leaderboard.router)
|
||
|
||
# Public embed Q&A status polling (/api/public/embed/*)
|
||
app.include_router(public_embed.router)
|
||
app.include_router(public_config.router)
|
||
app.include_router(html_page_favorites.router)
|
||
|
||
# Product knowledge base API is mounted at /api/knowledge
|
||
app.include_router(knowledge.router)
|
||
# Anonymous, section-level vector search over every published Wiki base
|
||
# except the system conversation-deposit base.
|
||
app.include_router(public_knowledge_vector_search.router)
|
||
|
||
# External knowledge ingest (publish md artifacts) at /api/knowledge-ingest
|
||
app.include_router(knowledge_ingest.router)
|
||
|
||
# System settings API is mounted at /api/system-settings
|
||
app.include_router(system_settings.router)
|
||
|
||
# Admin LLM metrics audit + export
|
||
app.include_router(llm_metrics.router)
|
||
|
||
# Admin window into other users' threads / messages / memory
|
||
app.include_router(admin_users.router)
|
||
app.include_router(admin_users.self_router)
|
||
|
||
app.include_router(admin_active_runs.router)
|
||
|
||
# Admin tool / skill call metrics audit + export
|
||
app.include_router(tool_metrics.router)
|
||
|
||
# Channels API is mounted at /api/channels
|
||
app.include_router(channels.router)
|
||
|
||
# Admin checkpoint migration API: SQLite legacy checkpoints -> Postgres.
|
||
app.include_router(checkpoint_migration.router)
|
||
|
||
# Assistants compatibility API (LangGraph Platform stub)
|
||
app.include_router(assistants_compat.router)
|
||
|
||
# Auth API is mounted at /api/v1/auth
|
||
app.include_router(auth.router)
|
||
|
||
# Feedback API is mounted at /api/threads/{thread_id}/runs/{run_id}/feedback
|
||
app.include_router(feedback.router)
|
||
|
||
# Notifications API is mounted at /api/notifications
|
||
app.include_router(notifications.router)
|
||
|
||
# Tags API is mounted at /api/tags
|
||
app.include_router(tags.router)
|
||
|
||
# 岗位 (position) management API is mounted at /api/positions
|
||
app.include_router(positions.router)
|
||
|
||
# Thread Runs API (LangGraph Platform-compatible runs lifecycle)
|
||
app.include_router(thread_runs.router)
|
||
|
||
app.include_router(thread_shares.router)
|
||
# Public, unauthenticated read-only share snapshots (/api/public/shares/*).
|
||
app.include_router(thread_shares.public_router)
|
||
|
||
# ?????????import ???? / view ????????
|
||
app.include_router(roundtable_draft_shares.router)
|
||
# Public, unauthenticated read-only roundtable result (/api/public/roundtable-shares/*).
|
||
app.include_router(roundtable_draft_shares.public_router)
|
||
|
||
# Stateless Runs API (stream/wait without a pre-existing thread)
|
||
app.include_router(runs.router)
|
||
|
||
# 开放问答接口 /api/open/chat —— 无需登录,指定模型 + 智能体做问答,可关闭思考模式
|
||
app.include_router(open_chat.router)
|
||
|
||
# 开放接口 /api/open/3qfx/ask —— 无需登录,按 taskId 直接发起 3qfx 第一次问答(后台运行,落成真实会话)
|
||
app.include_router(open_3qfx.router)
|
||
|
||
# --- ?????????? /page/workspace/roundtable/planning?---
|
||
# ??? app/gateway/routers/{multi_agent,intent,recommend}.py?
|
||
# ??? frontend-web/docs/multi-agent-backend-dev.md
|
||
app.include_router(multi_agent.router) # Step 2?init(SSE) + run/stream(leader|special)
|
||
app.include_router(intent.router) # Step 1?/api/intent/init + /stream
|
||
app.include_router(recommend.router) # Step 1?2?/api/recommend/stream?? init?
|
||
# Step 3 ????/?????? artifacts.router?/api/threads/{id}/artifacts?
|
||
|
||
# Health probes are always unauthenticated (see
|
||
# auth_middleware.PUBLIC_HEALTH_PATHS). Support the bare Gateway path,
|
||
# API base URL, and LangGraph proxy base URL for deployment monitors.
|
||
@app.api_route("/health", methods=["GET", "HEAD"], tags=["health"])
|
||
@app.api_route("/api/health", methods=["GET", "HEAD"], tags=["health"])
|
||
@app.api_route("/api/langgraph/health", methods=["GET", "HEAD"], tags=["health"])
|
||
async def health_check() -> dict:
|
||
"""Health check endpoint.
|
||
|
||
Returns:
|
||
Service health status information.
|
||
"""
|
||
return {"status": "healthy", "service": "deer-flow-gateway"}
|
||
|
||
# WeKnora's embedded frontend uses root-relative routes such as /platform,
|
||
# /assets, /locales and /api/v1/*. Register that proxy last so it stays a
|
||
# true fallback and cannot intercept DeerFlow APIs.
|
||
# Register the public-prefix routes before the unprefixed static fallback.
|
||
# Otherwise proxy_router's /{proxied_path:path} route consumes every
|
||
# /deerflow/* request before the prefixed platform/assets routes can match.
|
||
app.include_router(llmwiki.proxy_router, prefix="/deerflow")
|
||
app.include_router(llmwiki.proxy_router)
|
||
|
||
return app
|
||
|
||
|
||
# Create app instance for uvicorn
|
||
app = create_app()
|