mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-07-27 16:37:55 +00:00
* feat(channels): add GitHub event-driven agents (#3754) Add a webhook-driven GitHub channel with fail-closed webhook routing, deterministic per-agent PR/issue threads, mention-gated trigger fan-out, GitHub App token injection for sandboxed gh/git commands, and backend/AGENTS.md documentation. * fix(llm-middleware): classify bare IndexError as transient Upstream chat providers occasionally return 200 OK with an empty generations list (observed against Volces "coding" on ark.cn-beijing.volces.com). When that happens, langchain_core.language_models.chat_models.ainvoke raises ``IndexError: list index out of range`` at ``llm_result.generations[0][0].message`` and kills the run. Treat a bare IndexError reaching the middleware as a transient upstream-payload glitch and route it through the existing retry/backoff path instead of failing the whole agent run. The retry budget and backoff schedule are unchanged. Adds three regression tests covering the classifier and both the recover-on-retry and exhausted-retries paths. * fix(runtime): ignore stale LLM fallback markers from prior runs When a run on a thread ends with the LLM-error-handling middleware emitting a `deerflow_error_fallback`-marked AIMessage (e.g. after the IndexError empty-generations classification fix lands), that message is persisted to the thread's checkpoint as part of the messages channel. LangGraph replays the full message history in `stream_mode="values"` chunks, so every subsequent run on the same thread re-streams the stale fallback marker — and the worker's chunk scanner faithfully picks it up, flipping `RunStatus.success` to `RunStatus.error` for runs that themselves had no LLM failure at all. Snapshot the set of pre-existing message ids from the pre-run checkpoint and thread it through `_extract_llm_error_fallback_message` / `_try_extract_from_message` as a filter. Markers on history messages are ignored; markers on fresh messages produced during this run still trip the error path. Falls back to an empty set when the checkpointer is absent or the snapshot can't be captured, preserving the prior behavior on first-run / no-state paths. Adds unit tests for the new filter (helper-level and `_collect_pre_existing_message_ids`) plus an integration test exercising the full `run_agent` path with a stale history checkpointer. * fix(channels): make github channel fire-and-forget to avoid httpx.ReadTimeout on long runs GitHub agent runs (clone -> edit -> test -> push -> PR) routinely exceed the langgraph_sdk default 300s read deadline. The manager's runs.wait call kept an HTTP stream open for the entire run lifetime, so the long run blew up with httpx.ReadTimeout and the outer except branch then released the dedupe key and emitted a false 'internal error' outbound. The GitHub channel's outbound send is log-only by design: agents post to the issue/PR via the gh CLI in the sandbox when they choose to comment or create a PR. There is nothing for the manager to ferry back, so the long-poll was pure overhead. This change adds ChannelRunPolicy.fire_and_forget (default False) and sets it True for the github channel. When fire_and_forget is True, _handle_chat dispatches via client.runs.create (short POST, returns once the run is pending) instead of client.runs.wait, and skips the response-extraction + outbound-publish block. ConflictError on a busy thread still trips the standard THREAD_BUSY_MESSAGE path so behavior on the busy case is preserved for any future non-github fire-and-forget channel. Other (non-github) channels are unchanged: their policy defaults fire_and_forget=False and they continue to dispatch via runs.wait. Adds 6 regression tests in tests/test_channels.py::TestGithubFireAndForget: - Default ChannelRunPolicy.fire_and_forget is False. - The github policy registers fire_and_forget=True. - github inbound calls runs.create, not runs.wait, with the right kwargs. - github inbound publishes no outbound on success. - ConflictError from runs.create still emits THREAD_BUSY_MESSAGE. - Non-github channels (slack) still dispatch via runs.wait. * test(lead-agent): accept user_id kwarg in skill-policy test stubs The two GitHub-channel tests added in #3754 stubbed _load_enabled_skills_for_tool_policy with a lambda that only accepted `available_skills` and `app_config`, but the real function (and its call site in agent.py) also passes `user_id`. This raised TypeError on every run, failing backend-unit-tests. Add `user_id=None` to match the three sibling stubs in the same file. * refactor(gateway): disambiguate context-key set names The two frozensets _INTERNAL_ONLY_CONTEXT_KEYS and _CONTEXT_ONLY_KEYS shared a confusable "CONTEXT_ONLY" token in different orders, and the first broke the _CONTEXT_<X>_KEYS pattern of its sibling _CONTEXT_CONFIGURABLE_KEYS. Rename to make the distinct axes explicit: _CONTEXT_INTERNAL_CALLER_KEYS - WHO: internal callers (scheduler) only _CONTEXT_RUNTIME_ONLY_KEYS - WHERE: runtime context only, never configurable Pure rename, no behavior change.
483 lines
18 KiB
Python
483 lines
18 KiB
Python
"""CRUD API for custom agents."""
|
|
|
|
import asyncio
|
|
import logging
|
|
import re
|
|
import shutil
|
|
|
|
import yaml
|
|
from fastapi import APIRouter, HTTPException
|
|
from pydantic import BaseModel, Field
|
|
|
|
from deerflow.config.agents_api_config import get_agents_api_config
|
|
from deerflow.config.agents_config import AgentConfig, list_custom_agents, load_agent_config, load_agent_soul, preserve_non_managed_fields
|
|
from deerflow.config.paths import get_paths
|
|
from deerflow.runtime.user_context import get_effective_user_id
|
|
|
|
logger = logging.getLogger(__name__)
|
|
router = APIRouter(prefix="/api", tags=["agents"])
|
|
|
|
AGENT_NAME_PATTERN = re.compile(r"^[A-Za-z0-9-]+$")
|
|
|
|
|
|
class AgentResponse(BaseModel):
|
|
"""Response model for a custom agent."""
|
|
|
|
name: str = Field(..., description="Agent name (hyphen-case)")
|
|
description: str = Field(default="", description="Agent description")
|
|
model: str | None = Field(default=None, description="Optional model override")
|
|
tool_groups: list[str] | None = Field(default=None, description="Optional tool group whitelist")
|
|
skills: list[str] | None = Field(default=None, description="Optional skill whitelist (None=all, []=none)")
|
|
soul: str | None = Field(default=None, description="SOUL.md content")
|
|
|
|
|
|
class AgentsListResponse(BaseModel):
|
|
"""Response model for listing all custom agents."""
|
|
|
|
agents: list[AgentResponse]
|
|
|
|
|
|
class AgentCreateRequest(BaseModel):
|
|
"""Request body for creating a custom agent."""
|
|
|
|
name: str = Field(..., description="Agent name (must match ^[A-Za-z0-9-]+$, stored as lowercase)")
|
|
description: str = Field(default="", description="Agent description")
|
|
model: str | None = Field(default=None, description="Optional model override")
|
|
tool_groups: list[str] | None = Field(default=None, description="Optional tool group whitelist")
|
|
skills: list[str] | None = Field(default=None, description="Optional skill whitelist (None=all enabled, []=none)")
|
|
soul: str = Field(default="", description="SOUL.md content — agent personality and behavioral guardrails")
|
|
|
|
|
|
class AgentUpdateRequest(BaseModel):
|
|
"""Request body for updating a custom agent."""
|
|
|
|
description: str | None = Field(default=None, description="Updated description")
|
|
model: str | None = Field(default=None, description="Updated model override")
|
|
tool_groups: list[str] | None = Field(default=None, description="Updated tool group whitelist")
|
|
skills: list[str] | None = Field(default=None, description="Updated skill whitelist (None=all, []=none)")
|
|
soul: str | None = Field(default=None, description="Updated SOUL.md content")
|
|
|
|
|
|
def _validate_agent_name(name: str) -> None:
|
|
"""Validate agent name against allowed pattern.
|
|
|
|
Args:
|
|
name: The agent name to validate.
|
|
|
|
Raises:
|
|
HTTPException: 422 if the name is invalid.
|
|
"""
|
|
if not AGENT_NAME_PATTERN.match(name):
|
|
raise HTTPException(
|
|
status_code=422,
|
|
detail=f"Invalid agent name '{name}'. Must match ^[A-Za-z0-9-]+$ (letters, digits, and hyphens only).",
|
|
)
|
|
|
|
|
|
def _normalize_agent_name(name: str) -> str:
|
|
"""Normalize agent name to lowercase for filesystem storage."""
|
|
return name.lower()
|
|
|
|
|
|
def _require_agents_api_enabled() -> None:
|
|
"""Reject access unless the custom-agent management API is explicitly enabled."""
|
|
if not get_agents_api_config().enabled:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail=("Custom-agent management API is disabled. Set agents_api.enabled=true to expose agent and user-profile routes over HTTP."),
|
|
)
|
|
|
|
|
|
def _agent_config_to_response(agent_cfg: AgentConfig, include_soul: bool = False, *, user_id: str | None = None) -> AgentResponse:
|
|
"""Convert AgentConfig to AgentResponse."""
|
|
soul: str | None = None
|
|
if include_soul:
|
|
soul = load_agent_soul(agent_cfg.name, user_id=user_id) or ""
|
|
|
|
return AgentResponse(
|
|
name=agent_cfg.name,
|
|
description=agent_cfg.description,
|
|
model=agent_cfg.model,
|
|
tool_groups=agent_cfg.tool_groups,
|
|
skills=agent_cfg.skills,
|
|
soul=soul,
|
|
)
|
|
|
|
|
|
@router.get(
|
|
"/agents",
|
|
response_model=AgentsListResponse,
|
|
summary="List Custom Agents",
|
|
description="List all custom agents available in the agents directory, including their soul content.",
|
|
)
|
|
async def list_agents() -> AgentsListResponse:
|
|
"""List all custom agents.
|
|
|
|
Returns:
|
|
List of all custom agents with their metadata and soul content.
|
|
"""
|
|
_require_agents_api_enabled()
|
|
|
|
user_id = get_effective_user_id()
|
|
try:
|
|
agents = list_custom_agents(user_id=user_id)
|
|
return AgentsListResponse(agents=[_agent_config_to_response(a, include_soul=True, user_id=user_id) for a in agents])
|
|
except Exception as e:
|
|
logger.error(f"Failed to list agents: {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=f"Failed to list agents: {str(e)}")
|
|
|
|
|
|
@router.get(
|
|
"/agents/check",
|
|
summary="Check Agent Name",
|
|
description="Validate an agent name and check if it is available (case-insensitive).",
|
|
)
|
|
async def check_agent_name(name: str) -> dict:
|
|
"""Check whether an agent name is valid and not yet taken.
|
|
|
|
Args:
|
|
name: The agent name to check.
|
|
|
|
Returns:
|
|
``{"available": true/false, "name": "<normalized>"}``
|
|
|
|
Raises:
|
|
HTTPException: 422 if the name is invalid.
|
|
"""
|
|
_require_agents_api_enabled()
|
|
_validate_agent_name(name)
|
|
normalized = _normalize_agent_name(name)
|
|
user_id = get_effective_user_id()
|
|
paths = get_paths()
|
|
# Treat the name as taken if either the per-user path or the legacy shared
|
|
# path holds an agent — picking a name that collides with an unmigrated
|
|
# legacy agent would shadow the legacy entry once migration runs.
|
|
available = not paths.user_agent_dir(user_id, normalized).exists() and not paths.agent_dir(normalized).exists()
|
|
return {"available": available, "name": normalized}
|
|
|
|
|
|
@router.get(
|
|
"/agents/{name}",
|
|
response_model=AgentResponse,
|
|
summary="Get Custom Agent",
|
|
description="Retrieve details and SOUL.md content for a specific custom agent.",
|
|
)
|
|
async def get_agent(name: str) -> AgentResponse:
|
|
"""Get a specific custom agent by name.
|
|
|
|
Args:
|
|
name: The agent name.
|
|
|
|
Returns:
|
|
Agent details including SOUL.md content.
|
|
|
|
Raises:
|
|
HTTPException: 404 if agent not found.
|
|
"""
|
|
_require_agents_api_enabled()
|
|
_validate_agent_name(name)
|
|
name = _normalize_agent_name(name)
|
|
user_id = get_effective_user_id()
|
|
|
|
try:
|
|
agent_cfg = load_agent_config(name, user_id=user_id)
|
|
return _agent_config_to_response(agent_cfg, include_soul=True, user_id=user_id)
|
|
except FileNotFoundError:
|
|
raise HTTPException(status_code=404, detail=f"Agent '{name}' not found")
|
|
except Exception as e:
|
|
logger.error(f"Failed to get agent '{name}': {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=f"Failed to get agent: {str(e)}")
|
|
|
|
|
|
@router.post(
|
|
"/agents",
|
|
response_model=AgentResponse,
|
|
status_code=201,
|
|
summary="Create Custom Agent",
|
|
description="Create a new custom agent with its config and SOUL.md.",
|
|
)
|
|
async def create_agent_endpoint(request: AgentCreateRequest) -> AgentResponse:
|
|
"""Create a new custom agent.
|
|
|
|
Args:
|
|
request: The agent creation request.
|
|
|
|
Returns:
|
|
The created agent details.
|
|
|
|
Raises:
|
|
HTTPException: 409 if agent already exists, 422 if name is invalid.
|
|
"""
|
|
_require_agents_api_enabled()
|
|
_validate_agent_name(request.name)
|
|
normalized_name = _normalize_agent_name(request.name)
|
|
user_id = get_effective_user_id()
|
|
paths = get_paths()
|
|
|
|
def _create_agent() -> AgentResponse | None:
|
|
# Worker thread: base-dir resolution, existence checks, directory/file
|
|
# creation, read-back, and failure cleanup are all blocking filesystem
|
|
# IO that must stay off the event loop.
|
|
agent_dir = paths.user_agent_dir(user_id, normalized_name)
|
|
legacy_dir = paths.agent_dir(normalized_name)
|
|
|
|
if legacy_dir.exists():
|
|
return None # signals 409 to the caller
|
|
|
|
try:
|
|
try:
|
|
agent_dir.mkdir(parents=True, exist_ok=False)
|
|
except FileExistsError:
|
|
return None # signals 409 to the caller
|
|
# Write config.yaml
|
|
config_data: dict = {"name": normalized_name}
|
|
if request.description:
|
|
config_data["description"] = request.description
|
|
if request.model is not None:
|
|
config_data["model"] = request.model
|
|
if request.tool_groups is not None:
|
|
config_data["tool_groups"] = request.tool_groups
|
|
if request.skills is not None:
|
|
config_data["skills"] = request.skills
|
|
|
|
config_file = agent_dir / "config.yaml"
|
|
with open(config_file, "w", encoding="utf-8") as f:
|
|
yaml.dump(config_data, f, default_flow_style=False, allow_unicode=True)
|
|
|
|
# Write SOUL.md
|
|
soul_file = agent_dir / "SOUL.md"
|
|
soul_file.write_text(request.soul, encoding="utf-8")
|
|
|
|
logger.info(f"Created agent '{normalized_name}' at {agent_dir}")
|
|
|
|
agent_cfg = load_agent_config(normalized_name, user_id=user_id)
|
|
return _agent_config_to_response(agent_cfg, include_soul=True, user_id=user_id)
|
|
except Exception:
|
|
# Clean up partial state on failure before surfacing the error.
|
|
if agent_dir.exists():
|
|
shutil.rmtree(agent_dir)
|
|
raise
|
|
|
|
try:
|
|
response = await asyncio.to_thread(_create_agent)
|
|
except Exception as e:
|
|
logger.error(f"Failed to create agent '{request.name}': {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=f"Failed to create agent: {str(e)}")
|
|
|
|
if response is None:
|
|
raise HTTPException(status_code=409, detail=f"Agent '{normalized_name}' already exists")
|
|
|
|
return response
|
|
|
|
|
|
@router.put(
|
|
"/agents/{name}",
|
|
response_model=AgentResponse,
|
|
summary="Update Custom Agent",
|
|
description="Update an existing custom agent's config and/or SOUL.md.",
|
|
)
|
|
async def update_agent(name: str, request: AgentUpdateRequest) -> AgentResponse:
|
|
"""Update an existing custom agent.
|
|
|
|
Args:
|
|
name: The agent name.
|
|
request: The update request (all fields optional).
|
|
|
|
Returns:
|
|
The updated agent details.
|
|
|
|
Raises:
|
|
HTTPException: 404 if agent not found.
|
|
"""
|
|
_require_agents_api_enabled()
|
|
_validate_agent_name(name)
|
|
name = _normalize_agent_name(name)
|
|
user_id = get_effective_user_id()
|
|
|
|
try:
|
|
agent_cfg = load_agent_config(name, user_id=user_id)
|
|
except FileNotFoundError:
|
|
raise HTTPException(status_code=404, detail=f"Agent '{name}' not found")
|
|
|
|
paths = get_paths()
|
|
agent_dir = paths.user_agent_dir(user_id, name)
|
|
if not agent_dir.exists() and paths.agent_dir(name).exists():
|
|
raise HTTPException(
|
|
status_code=409,
|
|
detail=(f"Agent '{name}' only exists in the legacy shared layout and is not scoped to a user. Run scripts/migrate_user_isolation.py to move legacy agents into the per-user layout before updating."),
|
|
)
|
|
|
|
try:
|
|
# Update config if any config fields changed
|
|
# Use model_fields_set to distinguish "field omitted" from "explicitly set to null".
|
|
# This is critical for skills where None means "inherit all" (not "don't change").
|
|
fields_set = request.model_fields_set
|
|
config_changed = bool(fields_set & {"description", "model", "tool_groups", "skills"})
|
|
|
|
if config_changed:
|
|
updated: dict = {
|
|
"name": agent_cfg.name,
|
|
"description": request.description if "description" in fields_set else agent_cfg.description,
|
|
}
|
|
new_model = request.model if "model" in fields_set else agent_cfg.model
|
|
if new_model is not None:
|
|
updated["model"] = new_model
|
|
|
|
new_tool_groups = request.tool_groups if "tool_groups" in fields_set else agent_cfg.tool_groups
|
|
if new_tool_groups is not None:
|
|
updated["tool_groups"] = new_tool_groups
|
|
|
|
# skills: None = inherit all, [] = no skills, ["a","b"] = whitelist
|
|
if "skills" in fields_set:
|
|
new_skills = request.skills
|
|
else:
|
|
new_skills = agent_cfg.skills
|
|
if new_skills is not None:
|
|
updated["skills"] = new_skills
|
|
|
|
# Carry forward every top-level AgentConfig field this route does
|
|
# not manage (currently ``github:``, plus any future field added
|
|
# to :class:`AgentConfig`). The harness ``update_agent`` tool uses
|
|
# the same helper, so an operator editing the agent description
|
|
# from the Web UI does not silently strip a hand-authored
|
|
# ``github:`` binding — which would otherwise leave the next
|
|
# webhook delivery unable to find the agent in the registry and
|
|
# silently no-op.
|
|
for key, value in preserve_non_managed_fields(agent_cfg).items():
|
|
updated.setdefault(key, value)
|
|
|
|
config_file = agent_dir / "config.yaml"
|
|
with open(config_file, "w", encoding="utf-8") as f:
|
|
yaml.dump(updated, f, default_flow_style=False, allow_unicode=True)
|
|
|
|
# Update SOUL.md if provided
|
|
if request.soul is not None:
|
|
soul_path = agent_dir / "SOUL.md"
|
|
soul_path.write_text(request.soul, encoding="utf-8")
|
|
|
|
logger.info(f"Updated agent '{name}'")
|
|
|
|
refreshed_cfg = load_agent_config(name, user_id=user_id)
|
|
return _agent_config_to_response(refreshed_cfg, include_soul=True, user_id=user_id)
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to update agent '{name}': {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=f"Failed to update agent: {str(e)}")
|
|
|
|
|
|
class UserProfileResponse(BaseModel):
|
|
"""Response model for the global user profile (USER.md)."""
|
|
|
|
content: str | None = Field(default=None, description="USER.md content, or null if not yet created")
|
|
|
|
|
|
class UserProfileUpdateRequest(BaseModel):
|
|
"""Request body for setting the global user profile."""
|
|
|
|
content: str = Field(default="", description="USER.md content — describes the user's background and preferences")
|
|
|
|
|
|
@router.get(
|
|
"/user-profile",
|
|
response_model=UserProfileResponse,
|
|
summary="Get User Profile",
|
|
description="Read the global USER.md file that is injected into all custom agents.",
|
|
)
|
|
async def get_user_profile() -> UserProfileResponse:
|
|
"""Return the current USER.md content.
|
|
|
|
Returns:
|
|
UserProfileResponse with content=None if USER.md does not exist yet.
|
|
"""
|
|
_require_agents_api_enabled()
|
|
|
|
try:
|
|
user_md_path = get_paths().user_md_file
|
|
if not user_md_path.exists():
|
|
return UserProfileResponse(content=None)
|
|
raw = user_md_path.read_text(encoding="utf-8").strip()
|
|
return UserProfileResponse(content=raw or None)
|
|
except Exception as e:
|
|
logger.error(f"Failed to read user profile: {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=f"Failed to read user profile: {str(e)}")
|
|
|
|
|
|
@router.put(
|
|
"/user-profile",
|
|
response_model=UserProfileResponse,
|
|
summary="Update User Profile",
|
|
description="Write the global USER.md file that is injected into all custom agents.",
|
|
)
|
|
async def update_user_profile(request: UserProfileUpdateRequest) -> UserProfileResponse:
|
|
"""Create or overwrite the global USER.md.
|
|
|
|
Args:
|
|
request: The update request with the new USER.md content.
|
|
|
|
Returns:
|
|
UserProfileResponse with the saved content.
|
|
"""
|
|
_require_agents_api_enabled()
|
|
|
|
try:
|
|
paths = get_paths()
|
|
paths.base_dir.mkdir(parents=True, exist_ok=True)
|
|
paths.user_md_file.write_text(request.content, encoding="utf-8")
|
|
logger.info(f"Updated USER.md at {paths.user_md_file}")
|
|
return UserProfileResponse(content=request.content or None)
|
|
except Exception as e:
|
|
logger.error(f"Failed to update user profile: {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=f"Failed to update user profile: {str(e)}")
|
|
|
|
|
|
@router.delete(
|
|
"/agents/{name}",
|
|
status_code=204,
|
|
summary="Delete Custom Agent",
|
|
description="Delete a custom agent and all its files (config, SOUL.md, memory).",
|
|
)
|
|
async def delete_agent(name: str) -> None:
|
|
"""Delete a custom agent.
|
|
|
|
Args:
|
|
name: The agent name.
|
|
|
|
Raises:
|
|
HTTPException: 404 if no per-user copy exists; 409 if only a legacy
|
|
shared copy exists (suggesting the migration script).
|
|
"""
|
|
_require_agents_api_enabled()
|
|
_validate_agent_name(name)
|
|
name = _normalize_agent_name(name)
|
|
user_id = get_effective_user_id()
|
|
paths = get_paths()
|
|
|
|
def _remove_agent_dir() -> tuple[str, str]:
|
|
# Runs in a worker thread: resolving the base dir, probing the directory
|
|
# (`exists`), and removing it (`rmtree`) are all blocking filesystem IO
|
|
# that must stay off the event loop.
|
|
agent_dir = paths.user_agent_dir(user_id, name)
|
|
if not agent_dir.exists():
|
|
outcome = "legacy" if paths.agent_dir(name).exists() else "missing"
|
|
return outcome, str(agent_dir)
|
|
shutil.rmtree(agent_dir)
|
|
return "deleted", str(agent_dir)
|
|
|
|
try:
|
|
outcome, agent_dir = await asyncio.to_thread(_remove_agent_dir)
|
|
except Exception as e:
|
|
logger.error(f"Failed to delete agent '{name}': {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=f"Failed to delete agent: {str(e)}")
|
|
|
|
if outcome == "legacy":
|
|
raise HTTPException(
|
|
status_code=409,
|
|
detail=(f"Agent '{name}' only exists in the legacy shared layout and is not scoped to a user. Run scripts/migrate_user_isolation.py to move legacy agents into the per-user layout before deleting."),
|
|
)
|
|
if outcome == "missing":
|
|
raise HTTPException(status_code=404, detail=f"Agent '{name}' not found")
|
|
|
|
logger.info(f"Deleted agent '{name}' from {agent_dir}")
|