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()