deerflow-code/offline-backend-20260512/backend/app/gateway/routers/workflow_resources.py
2026-09-07 18:24:55 +08:00

248 lines
10 KiB
Python

"""Workflow node-type and resource catalog APIs (phase 1)."""
from __future__ import annotations
from typing import Any
from fastapi import APIRouter, HTTPException, Request
from app.gateway.deps import get_current_user, get_local_provider
from app.gateway.routers._workflow_planner_seed import WORKFLOW_PLANNER_AGENT_ID
from deerflow.config.app_config import get_app_config
from deerflow.persistence.workflows import WorkflowStore
from deerflow.workflows.schemas import ALL_NODE_TYPES
router = APIRouter(prefix="/api/workflows", tags=["workflow-resources"])
_NODE_META: dict[str, dict[str, str]] = {
"start": {"label": "开始", "category": "control"},
"output": {"label": "输出", "category": "control"},
"agent": {"label": "智能体", "category": "ai"},
"skill": {"label": "技能", "category": "ai"},
"http": {"label": "HTTP", "category": "data"},
"sql_read": {"label": "SQL 只读", "category": "data"},
"code": {"label": "代码", "category": "data"},
"transform": {"label": "模板/转换", "category": "data"},
"evidence_normalizer": {"label": "证据归一化", "category": "data"},
"deep_research_write": {"label": "深度研究报告", "category": "ai"},
"condition": {"label": "条件", "category": "control"},
"merge": {"label": "合并", "category": "control"},
"human_input": {"label": "人工介入", "category": "control"},
"subworkflow": {"label": "子工作流", "category": "control"},
"loop": {"label": "循环", "category": "control"},
}
def _require_enabled() -> None:
if not get_app_config().workflows.enabled:
raise HTTPException(status_code=503, detail="Workflow Studio is disabled")
def _get_store(request: Request) -> WorkflowStore:
store = getattr(request.app.state, "workflow_store", None)
if store is None:
raise HTTPException(status_code=503, detail="Workflow store not available")
return store
async def _resolve_creator_names(user_ids: set[str]) -> dict[str, str]:
"""Resolve resource-owner ids to short display names, best-effort.
The workflow catalog is a small display projection over the agent and
skill stores. Those stores intentionally retain immutable owner ids for
access control; this resolver exposes a friendly creator name without
changing that ownership model. If an account has since disappeared, the
stable id remains a useful audit fallback.
"""
if not user_ids:
return {}
try:
provider = get_local_provider()
except Exception: # noqa: BLE001
return {user_id: user_id for user_id in user_ids}
names: dict[str, str] = {}
for user_id in user_ids:
try:
user = await provider.get_user(user_id)
except Exception: # noqa: BLE001
user = None
email = getattr(user, "email", None) if user is not None else None
if isinstance(email, str) and email:
names[user_id] = email.split("@", 1)[0] or user_id
else:
names[user_id] = user_id
return names
def _creator_label(owner_id: object, creator_names: dict[str, str]) -> str:
"""Return the catalog label for a user-owned or built-in resource."""
if owner_id is None or not str(owner_id).strip():
return "系统内置"
owner_key = str(owner_id)
return creator_names.get(owner_key, owner_key)
@router.get("/node-types")
async def list_node_types(request: Request) -> dict[str, Any]:
_require_enabled()
await get_current_user(request)
items = []
for node_type in ALL_NODE_TYPES:
meta = _NODE_META.get(node_type, {"label": node_type, "category": "other"})
items.append({"type": node_type, **meta})
return {"nodeTypes": items}
@router.get("/resources/agents")
async def list_agent_resources(request: Request) -> dict[str, Any]:
_require_enabled()
user_id = await get_current_user(request)
agent_store = getattr(request.app.state, "agent_store", None)
agents: list[dict[str, Any]] = []
if agent_store is not None:
try:
if hasattr(agent_store, "list_visible"):
rows = await agent_store.list_visible(user_id)
elif hasattr(agent_store, "list_all"):
rows = await agent_store.list_all()
else:
rows = []
# This is a control-plane agent. It plans candidates through the
# planning endpoint and must never appear as a draggable worker.
rows = [
row
for row in rows or []
if isinstance(row, dict)
and str(row.get("id") or row.get("agent_id") or "") != WORKFLOW_PLANNER_AGENT_ID
]
creator_names = await _resolve_creator_names({str(row["user_id"]) for row in rows or [] if row.get("user_id") is not None})
agent_ids = [str(row.get("id") or row.get("agent_id") or "") for row in rows or []]
extras_by_agent: dict[str, dict[str, Any]] = {}
if hasattr(agent_store, "get_extras_for"):
extras_by_agent = await agent_store.get_extras_for([agent_id for agent_id in agent_ids if agent_id])
keyword = request.query_params.get("keyword", "").strip().casefold()
for row in rows or []:
agent_id = str(row.get("id") or row.get("agent_id") or "")
name = str(row.get("name") or agent_id)
description = str(row.get("description") or "")
if keyword and keyword not in f"{agent_id} {name} {description}".casefold():
continue
cached_skills = (extras_by_agent.get(agent_id) or {}).get("skills") or []
agents.append(
{
"agentId": agent_id,
"name": name,
"description": description,
"skills": [str(skill) for skill in cached_skills],
"createdBy": _creator_label(row.get("user_id"), creator_names),
}
)
except Exception: # noqa: BLE001
agents = []
return {"agents": agents}
@router.get("/resources/skills")
async def list_skill_resources(request: Request) -> dict[str, Any]:
_require_enabled()
user_id = await get_current_user(request)
skills: list[dict[str, Any]] = []
try:
from deerflow.skills import get_or_new_skill_storage
skill_store = getattr(request.app.state, "skill_store", None)
ownership_by_name: dict[str, dict[str, Any]] = {}
if skill_store is not None and hasattr(skill_store, "list_visible"):
records = await skill_store.list_visible(user_id)
ownership_by_name = {str(record["name"]): record for record in records if record.get("name")}
creator_names = await _resolve_creator_names({str(record["owner_user_id"]) for record in ownership_by_name.values() if record.get("owner_user_id") is not None})
storage = get_or_new_skill_storage()
for skill in storage.load_skills(enabled_only=True) or []:
name = getattr(skill, "name", None)
if not name:
continue
record = ownership_by_name.get(str(name), {})
skills.append(
{
"skillId": name,
"name": name,
"description": getattr(skill, "description", "") or "",
# Workflow direct-skill execution remains opt-in. Older
# Skill definitions do not advertise this capability and
# are therefore correctly shown as agent-only resources.
"callable": bool(getattr(skill, "callable", False)),
"createdBy": _creator_label(record.get("owner_user_id"), creator_names),
}
)
except Exception: # noqa: BLE001
skills = []
return {"skills": skills}
@router.get("/resources/data-sources")
async def list_data_source_resources(request: Request) -> dict[str, Any]:
_require_enabled()
user_id = await get_current_user(request)
store = getattr(request.app.state, "workflow_data_source_store", None)
if store is None:
return {"dataSources": []}
rows = await store.list_sources(owner_id=user_id)
return {
"dataSources": [
{
"id": row["id"],
"dataSourceId": row["id"],
"name": row["name"],
"description": row.get("description") or "",
"kind": row.get("kind"),
"engine": row.get("driver") or "",
"baseUrl": row.get("base_url") or "",
"maskedTarget": row.get("masked_target"),
"enabled": row.get("enabled"),
"readOnly": row.get("kind") == "sql",
"allowedMethods": row.get("allowed_methods") or [],
"allowedTables": row.get("allowed_tables") or [],
}
for row in rows
if row.get("enabled", True)
]
}
@router.get("/resources/credentials")
async def list_credential_resources(request: Request) -> dict[str, Any]:
"""Read-only credential directory: ids and labels only.
HTTP credentials are ``http``-kind data sources; the encrypted header blob
is resolved exclusively by the run executor. This endpoint must never
return headers, tokens, DSNs or ciphertext — a 404-free catalog so the
canvas credential picker does not fail the whole resource panel.
"""
_require_enabled()
user_id = await get_current_user(request)
store = getattr(request.app.state, "workflow_data_source_store", None)
if store is None:
return {"credentials": []}
rows = await store.list_sources(owner_id=user_id, include_shared=True)
return {
"credentials": [
{
"credentialId": row["id"],
"name": row["name"],
"type": "http_headers",
"sourceKind": "http",
}
for row in rows
if str(row.get("kind") or "") == "http" and row.get("enabled", True)
]
}
@router.get("/resources/subworkflows")
async def list_subworkflow_resources(request: Request) -> dict[str, Any]:
_require_enabled()
await get_current_user(request)
store = _get_store(request)
return {"subworkflows": await store.list_published_summaries()}