mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-07-28 17:06:05 +00:00
* feat(agents): database-backed storage for custom agent definitions Add an agent_storage.backend switch (default file, behaviour-unchanged) with a db backend that stores each custom agent as a row in the shared SQL persistence layer, so a multi-instance deployment sees the same agents on every node (#4331, #4357). Introduces an AgentStore interface routing all read/write surfaces, an agents table + migration 0006, startup validation, and a file->db importer. Follows the thread_meta store / run_events backend-switch / 0003_scheduled_tasks migration patterns; no new dependency. * fix(agents): make db storage path production-ready (review round 1) Addresses review feedback on the db/sync agent-storage path: - sql.py: mirror the async engine's per-connection SQLite PRAGMAs on the sync engine (busy_timeout=30000, synchronous=NORMAL, foreign_keys=ON, WAL) so both engines behave identically against the shared DB; guard the engine cache with a lock (double-checked) so concurrent first-touch cannot build duplicate engines or register the connect listener twice. - routers/agents.py + routers/assistants_compat.py: offload the sync-store reads that ran on the event loop (list/get/check, update's pre-read + legacy guard + refresh, and assistants_compat's four list routes) via asyncio.to_thread — on db+postgres each was a network round trip stalling the loop. Writes were already offloaded. - file.py: translate the create() mkdir(exist_ok=False) race FileExistsError into AgentExistsError (router 409, matching SqlAgentStore's IntegrityError path); correct the _write docstring — per-file atomic replace, two commits sequential not transactional. Tests: sync-engine PRAGMA + engine-cache reuse assertions; file create-race -> AgentExistsError; strict Blockbuster anchor over the read endpoints so a regression back onto the loop fails CI. * fix(agents): address round-2 review on the db store path - update_agent tool: align the docstring/inline comment with FileAgentStore._write. Cross-field write atomicity is db-only; the file backend commits config then soul via two sequential os.replace (a crash between them can leave a fresh config.yaml beside a stale SOUL.md). The dropped partial-write *reporting* is an intentional tradeoff — the stage-then-replace safety is preserved (test_update_agent_soul_failure_does_not_replace_config still holds). - SqlAgentStore.update(): true upsert. Catch IntegrityError on the insert-on-missing branch, re-fetch and apply, so two concurrent first-time writes (e.g. two setup_agent handshakes) converge instead of surfacing a raw UNIQUE(user_id, name) violation as a 500. Symmetric with create(). - get_agent_store(): document the graph-subprocess config-resolution invariant (the except->file fallback is a genuine no-config path, not a mask for a misconfigured graph process) and pin it with two tests driving the real get_app_config() file resolution: db resolves from an on-disk config.yaml, file fallback when config is unresolvable. * test(agents): cover SqlAgentStore.update() write-race upsert recovery Mandatory-TDD test for the round-2 fix in 0680340a: two concurrent first-time update()s where the loser's insert hits UNIQUE(user_id, name). Deterministically forces the IntegrityError recovery path by making the first _row probe miss the committed winner, and asserts last-writer-wins instead of a surfaced 500.
95 lines
4.1 KiB
Python
95 lines
4.1 KiB
Python
"""Regression anchors: the custom-agent router must not block the event loop.
|
|
|
|
``app.gateway.routers.agents.create_agent_endpoint`` and ``delete_agent`` are
|
|
async route handlers that resolve the agent directory (``Paths.base_dir`` calls
|
|
``Path.resolve``), probe it (``Path.exists``), and create/remove it (``mkdir``,
|
|
config/SOUL writes, ``shutil.rmtree``) — all blocking IO. Both offload that work
|
|
via ``asyncio.to_thread``; if any of it regresses back onto the event loop, the
|
|
strict Blockbuster gate raises ``BlockingError`` and these tests fail.
|
|
|
|
Imports live at module scope so the one-time FastAPI app construction (which
|
|
reads files while building OpenAPI schemas) happens at collection time, not on
|
|
the event loop under test. Test-side path resolution is itself offloaded with
|
|
``asyncio.to_thread`` (matching ``test_uploads_middleware``) so only the
|
|
handlers' own filesystem access is exercised on the loop.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
from app.gateway.routers.agents import (
|
|
AgentCreateRequest,
|
|
check_agent_name,
|
|
create_agent_endpoint,
|
|
delete_agent,
|
|
get_agent,
|
|
list_agents,
|
|
)
|
|
from deerflow.config.agents_api_config import load_agents_api_config_from_dict
|
|
from deerflow.config.paths import get_paths
|
|
from deerflow.runtime.user_context import get_effective_user_id
|
|
|
|
pytestmark = pytest.mark.asyncio
|
|
|
|
|
|
async def test_create_agent_does_not_block_event_loop(tmp_path: Path, monkeypatch) -> None:
|
|
monkeypatch.setenv("DEER_FLOW_HOME", str(tmp_path))
|
|
monkeypatch.setattr("deerflow.config.paths._paths", None)
|
|
load_agents_api_config_from_dict({"enabled": True})
|
|
try:
|
|
response = await create_agent_endpoint(AgentCreateRequest(name="loop-make-agent", soul="You are a test agent."))
|
|
assert response is not None
|
|
|
|
user_id = get_effective_user_id()
|
|
# test-side check (resolution offloaded; not exercised on the loop)
|
|
agent_dir = await asyncio.to_thread(get_paths().user_agent_dir, user_id, "loop-make-agent")
|
|
assert await asyncio.to_thread((agent_dir / "config.yaml").exists)
|
|
finally:
|
|
load_agents_api_config_from_dict({})
|
|
|
|
|
|
async def test_delete_agent_does_not_block_event_loop(tmp_path: Path, monkeypatch) -> None:
|
|
monkeypatch.setenv("DEER_FLOW_HOME", str(tmp_path))
|
|
monkeypatch.setattr("deerflow.config.paths._paths", None)
|
|
load_agents_api_config_from_dict({"enabled": True})
|
|
try:
|
|
user_id = get_effective_user_id()
|
|
user_id = get_effective_user_id()
|
|
# test-side seeding (resolution offloaded; not exercised on the loop)
|
|
agent_dir = await asyncio.to_thread(get_paths().user_agent_dir, user_id, "loop-test-agent")
|
|
await asyncio.to_thread(agent_dir.mkdir, parents=True, exist_ok=True)
|
|
await asyncio.to_thread((agent_dir / "config.yaml").write_text, "name: loop-test-agent\n", encoding="utf-8")
|
|
|
|
await delete_agent("loop-test-agent")
|
|
|
|
assert not await asyncio.to_thread(agent_dir.exists)
|
|
finally:
|
|
load_agents_api_config_from_dict({})
|
|
|
|
|
|
async def test_read_endpoints_do_not_block_event_loop(tmp_path: Path, monkeypatch) -> None:
|
|
# list/get/check read through the sync agent store; on the db backend each is
|
|
# a DB round trip. They must offload via asyncio.to_thread, or the strict
|
|
# Blockbuster gate raises BlockingError here (finding: reads on the loop).
|
|
monkeypatch.setenv("DEER_FLOW_HOME", str(tmp_path))
|
|
monkeypatch.setattr("deerflow.config.paths._paths", None)
|
|
load_agents_api_config_from_dict({"enabled": True})
|
|
try:
|
|
await create_agent_endpoint(AgentCreateRequest(name="loop-read-agent", soul="You are a test agent."))
|
|
|
|
listed = await list_agents()
|
|
assert any(a.name == "loop-read-agent" for a in listed.agents)
|
|
|
|
got = await get_agent("loop-read-agent")
|
|
assert got.name == "loop-read-agent"
|
|
|
|
check = await check_agent_name("loop-read-agent")
|
|
assert check["available"] is False
|
|
assert (await check_agent_name("never-created-agent"))["available"] is True
|
|
finally:
|
|
load_agents_api_config_from_dict({})
|