mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-09-12 23:19:36 +00:00
* fix: bound MCP server bring-up timeouts and exclude externalized tool outputs from delivery verification Two related robustness fixes: 1. MCP server bring-up was unbounded. tool_call_timeout only covered session.call_tool(); tool discovery (subprocess spawn + initialize + tools/list) and persistent stdio session initialization could hang forever, blocking agent construction (and on the Gateway event loop, the whole process). Add a per-server session_init_timeout (default DEFAULT_MCP_SESSION_INIT_TIMEOUT = 60s, null disables) that bounds both discovery and pooled-session initialization. The session pool's existing cancellation handling tears down a session stuck mid-creation in its own task. 2. ToolOutputBudgetMiddleware externalizes oversized tool outputs into outputs/.tool-results/ (configurable tool_output.storage_subdir). The workspace-change scanner and run delivery verification counted those files as produced artifacts, so any run that externalized a tool output without also presenting a real artifact failed with "Artifact delivery incomplete". Exclude TOOL_RESULTS_DIRNAME via a shared constant (mirroring BROWSER_FRAMES_DIRNAME) and thread the configured storage_subdir through snapshot capture so both workspace-changes events and delivery verification stay clean. * review: enforce single-segment tool_output.storage_subdir; document discovery-timeout cleanup Address review feedback: 1. A custom tool_output.storage_subdir with a path separator (e.g. cache/tool-results) silently no-oped the workspace-scanner exclusion: os.walk yields one-segment dirnames, so a nested value never matched and its files were counted as produced artifacts again. ToolOutputConfig now validates storage_subdir as a single directory name (rejects separators, .., absolute, empty) with tests, so the exclusion is always sound. 2. The discovery-timeout path now documents why cancellation is safe, mirroring the session-init note: discovery runs inside the adapter's nested async context managers, and stdio_client's finally terminates the process tree (SIGTERM->SIGKILL on POSIX, process-tree on Windows), so a timed-out npx subprocess and its children are reaped rather than accumulating. * review: log session-init timeouts and align API response model default with runtime config Address second-round review feedback: 1. A session-init timeout raised TimeoutError without any log, unlike the discovery timeout which logs a WARNING. Wrap the bounded get_session in a try/except that logs the timeout (server name + seconds) and re-raises, so operators can diagnose tool-call failures caused by hung MCP sessions. 2. McpServerConfigResponse.session_init_timeout defaulted to None while McpServerConfig defaults to 60s: a server created via PUT /api/mcp/config without the field was persisted with null (no timeout) while the same server created in the config file got 60s. Align the response-model default to DEFAULT_MCP_SESSION_INIT_TIMEOUT so API-created and file-created servers behave the same; an explicit null still opts out. * review: narrow the discovery-timeout handler to the bounded wait_for path The except TimeoutError clause covered both the bounded wait_for branch and the bare discovery branch. With session_init_timeout opted out (None), a TimeoutError raised by discovery itself would hit the %.1f format with None: logging raises TypeError internally, the WARNING is silently dropped, and a --- Logging error --- traceback goes to stderr. Narrow the handler to wrap only the wait_for call, where the branch condition guarantees the timeout value is not None. A discovery-internal TimeoutError on the opted-out path now falls through to the generic failure handler and is reported as 'tool discovery failed' with exc_info. Covered by a regression test that asserts the skip is reported without any broken format.
739 lines
27 KiB
Python
739 lines
27 KiB
Python
"""Worker-level regression tests for the terminal run.delivery event (#4272 slice 1)."""
|
|
|
|
import asyncio
|
|
from types import SimpleNamespace
|
|
from unittest.mock import AsyncMock, MagicMock
|
|
from uuid import uuid4
|
|
|
|
import pytest
|
|
from langchain_core.messages import AIMessage, ToolMessage
|
|
from langgraph.types import Command
|
|
|
|
from deerflow.config.app_config import AppConfig
|
|
from deerflow.config.paths import Paths
|
|
from deerflow.config.sandbox_config import SandboxConfig
|
|
from deerflow.config.tool_output_config import ToolOutputConfig
|
|
from deerflow.runtime.events.store.memory import MemoryRunEventStore
|
|
from deerflow.runtime.runs.manager import RunManager
|
|
from deerflow.runtime.runs.schemas import RunStatus
|
|
from deerflow.runtime.runs.store.memory import MemoryRunStore
|
|
from deerflow.runtime.runs.worker import RunContext, _delivery_content_with_outputs, run_agent
|
|
from deerflow.runtime.user_context import get_effective_user_id
|
|
|
|
|
|
def _make_bridge():
|
|
return SimpleNamespace(publish=AsyncMock(), publish_end=AsyncMock(), cleanup=AsyncMock())
|
|
|
|
|
|
async def _delivery_events(store: MemoryRunEventStore, thread_id: str, run_id: str) -> list[dict]:
|
|
events = await store.list_events(thread_id, run_id)
|
|
return [e for e in events if e["event_type"] == "run.delivery"]
|
|
|
|
|
|
def test_delivery_verification_treats_presented_directory_as_covering_produced_files():
|
|
content = {
|
|
"presented": 1,
|
|
"paths": ["/mnt/user-data/outputs/site"],
|
|
"by_tool": {"present_files": ["/mnt/user-data/outputs/site"]},
|
|
}
|
|
|
|
delivery = _delivery_content_with_outputs(
|
|
content,
|
|
[
|
|
"/mnt/user-data/outputs/site/index.html",
|
|
"/mnt/user-data/outputs/site/assets/style.css",
|
|
],
|
|
)
|
|
|
|
assert delivery["matched_paths"] == [
|
|
"/mnt/user-data/outputs/site/index.html",
|
|
"/mnt/user-data/outputs/site/assets/style.css",
|
|
]
|
|
assert delivery["satisfied"] is True
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_event_records_present_files_paths_on_success():
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
|
|
class DummyAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
journal = config["context"]["__run_journal"]
|
|
ai = AIMessage(content="", tool_calls=[{"id": "call_1", "name": "present_files", "args": {}}])
|
|
journal._remember_current_run_tool_calls(ai, caller="lead_agent")
|
|
journal.on_tool_end(
|
|
Command(
|
|
update={
|
|
"artifacts": ["/mnt/user-data/outputs/report.md"],
|
|
"messages": [ToolMessage("Successfully presented files", tool_call_id="call_1")],
|
|
}
|
|
),
|
|
run_id=uuid4(),
|
|
)
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: DummyAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
await asyncio.sleep(0)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert len(delivery) == 1
|
|
assert delivery[0]["content"]["presented"] == 1
|
|
assert delivery[0]["content"]["paths"] == ["/mnt/user-data/outputs/report.md"]
|
|
assert delivery[0]["content"]["by_tool"] == {"present_files": ["/mnt/user-data/outputs/report.md"]}
|
|
fetched = await run_manager.get(record.run_id)
|
|
assert fetched.status == RunStatus.success
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_event_presented_zero_without_artifact_production():
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
|
|
class DummyAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: DummyAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
await asyncio.sleep(0)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert len(delivery) == 1
|
|
assert delivery[0]["content"] == {"presented": 0, "paths": [], "by_tool": {}}
|
|
fetched = await run_manager.get(record.run_id)
|
|
assert fetched.status == RunStatus.success
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_changed_outputs_succeed_when_a_produced_output_is_presented(monkeypatch):
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
monkeypatch.setattr(
|
|
"deerflow.runtime.runs.worker._produced_output_paths",
|
|
AsyncMock(return_value=["/mnt/user-data/outputs/report.md"]),
|
|
)
|
|
|
|
class DummyAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
journal = config["context"]["__run_journal"]
|
|
ai = AIMessage(content="", tool_calls=[{"id": "call_1", "name": "present_files", "args": {}}])
|
|
journal._remember_current_run_tool_calls(ai, caller="lead_agent")
|
|
journal.on_tool_end(
|
|
Command(
|
|
update={
|
|
"artifacts": ["/mnt/user-data/outputs/report.md"],
|
|
"messages": [ToolMessage("Successfully presented files", tool_call_id="call_1")],
|
|
}
|
|
),
|
|
run_id=uuid4(),
|
|
)
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: DummyAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert delivery[0]["content"] == {
|
|
"presented": 1,
|
|
"paths": ["/mnt/user-data/outputs/report.md"],
|
|
"by_tool": {"present_files": ["/mnt/user-data/outputs/report.md"]},
|
|
"verification": {
|
|
"source": "outputs_changed",
|
|
"requirement": "present_files_matches_produced_output",
|
|
},
|
|
"produced_paths": ["/mnt/user-data/outputs/report.md"],
|
|
"presented_paths": ["/mnt/user-data/outputs/report.md"],
|
|
"matched_paths": ["/mnt/user-data/outputs/report.md"],
|
|
"stage": "presented",
|
|
"satisfied": True,
|
|
}
|
|
assert record.status == RunStatus.success
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_changed_outputs_fail_closed_when_not_presented(monkeypatch):
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
monkeypatch.setattr(
|
|
"deerflow.runtime.runs.worker._produced_output_paths",
|
|
AsyncMock(return_value=["/mnt/user-data/outputs/report.md"]),
|
|
)
|
|
|
|
class ProseOnlyAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
yield {"messages": [AIMessage(content="SESSION SUMMARY")]}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: ProseOnlyAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert delivery[0]["content"] == {
|
|
"presented": 0,
|
|
"paths": [],
|
|
"by_tool": {},
|
|
"verification": {
|
|
"source": "outputs_changed",
|
|
"requirement": "present_files_matches_produced_output",
|
|
},
|
|
"produced_paths": ["/mnt/user-data/outputs/report.md"],
|
|
"presented_paths": [],
|
|
"matched_paths": [],
|
|
"stage": "not_started",
|
|
"satisfied": False,
|
|
}
|
|
assert record.status == RunStatus.error
|
|
assert record.error == "Artifact delivery incomplete: no produced output artifact was presented"
|
|
assert record.stop_reason is None
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_externalized_tool_results_do_not_trigger_delivery_verification(tmp_path, monkeypatch):
|
|
"""Oversized tool outputs externalized under outputs/.tool-results/ are
|
|
process feedback for the model, not deliverables: a run that only produced
|
|
those files must succeed without any present_files call."""
|
|
paths = Paths(base_dir=tmp_path)
|
|
monkeypatch.setattr("deerflow.workspace_changes.recorder.get_paths", lambda: paths)
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
|
|
class ExternalizingAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
# Simulates ToolOutputBudgetMiddleware persisting an oversized tool
|
|
# output mid-run (default storage_subdir is ".tool-results").
|
|
tool_results = paths.sandbox_outputs_dir("thread-1", user_id=get_effective_user_id()) / ".tool-results"
|
|
tool_results.mkdir(parents=True, exist_ok=True)
|
|
(tool_results / "bash-abcdef123456.log").write_text("x" * 20000, encoding="utf-8")
|
|
yield {"messages": [AIMessage(content="Here is the answer.")]}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: ExternalizingAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert len(delivery) == 1
|
|
assert delivery[0]["content"] == {"presented": 0, "paths": [], "by_tool": {}}
|
|
fetched = await run_manager.get(record.run_id)
|
|
assert fetched.status == RunStatus.success
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_custom_tool_output_storage_subdir_does_not_trigger_delivery_verification(tmp_path, monkeypatch):
|
|
"""A custom tool_output.storage_subdir is honoured by the exclusion, not
|
|
only the default .tool-results name."""
|
|
paths = Paths(base_dir=tmp_path)
|
|
monkeypatch.setattr("deerflow.workspace_changes.recorder.get_paths", lambda: paths)
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
app_config = AppConfig(sandbox=SandboxConfig(use="test"), tool_output=ToolOutputConfig(storage_subdir="tool-output-cache"))
|
|
|
|
class ExternalizingAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
cache = paths.sandbox_outputs_dir("thread-1", user_id=get_effective_user_id()) / "tool-output-cache"
|
|
cache.mkdir(parents=True, exist_ok=True)
|
|
(cache / "web_fetch-abcdef123456.log").write_text("y" * 20000, encoding="utf-8")
|
|
yield {"messages": [AIMessage(content="Here is the answer.")]}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store, app_config=app_config),
|
|
agent_factory=lambda *, config: ExternalizingAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
fetched = await run_manager.get(record.run_id)
|
|
assert fetched.status == RunStatus.success
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_changed_outputs_succeed_when_one_of_multiple_outputs_is_presented(monkeypatch):
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
monkeypatch.setattr(
|
|
"deerflow.runtime.runs.worker._produced_output_paths",
|
|
AsyncMock(
|
|
return_value=[
|
|
"/mnt/user-data/outputs/report.md",
|
|
"/mnt/user-data/outputs/appendix.md",
|
|
]
|
|
),
|
|
)
|
|
|
|
class PartiallyPresentingAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
journal = config["context"]["__run_journal"]
|
|
journal._remember_current_run_tool_calls(
|
|
AIMessage(content="", tool_calls=[{"id": "call_1", "name": "present_files", "args": {}}]),
|
|
caller="lead_agent",
|
|
)
|
|
journal.on_tool_end(
|
|
Command(
|
|
update={
|
|
"artifacts": ["/mnt/user-data/outputs/report.md"],
|
|
"messages": [ToolMessage("Successfully presented files", tool_call_id="call_1")],
|
|
}
|
|
),
|
|
run_id=uuid4(),
|
|
)
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: PartiallyPresentingAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = (await _delivery_events(store, "thread-1", record.run_id))[0]["content"]
|
|
assert delivery["stage"] == "presented"
|
|
assert delivery["presented_paths"] == ["/mnt/user-data/outputs/report.md"]
|
|
assert delivery["matched_paths"] == ["/mnt/user-data/outputs/report.md"]
|
|
assert delivery["satisfied"] is True
|
|
assert record.status == RunStatus.success
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_changed_outputs_fail_when_present_files_only_presents_an_unrelated_file(monkeypatch):
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
monkeypatch.setattr(
|
|
"deerflow.runtime.runs.worker._produced_output_paths",
|
|
AsyncMock(return_value=["/mnt/user-data/outputs/report.md"]),
|
|
)
|
|
|
|
class UnrelatedPresentingAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
journal = config["context"]["__run_journal"]
|
|
journal._remember_current_run_tool_calls(
|
|
AIMessage(content="", tool_calls=[{"id": "call_1", "name": "present_files", "args": {}}]),
|
|
caller="lead_agent",
|
|
)
|
|
journal.on_tool_end(
|
|
Command(
|
|
update={
|
|
"artifacts": ["/mnt/user-data/outputs/old-report.md"],
|
|
"messages": [ToolMessage("Successfully presented files", tool_call_id="call_1")],
|
|
}
|
|
),
|
|
run_id=uuid4(),
|
|
)
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: UnrelatedPresentingAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = (await _delivery_events(store, "thread-1", record.run_id))[0]["content"]
|
|
assert delivery["stage"] == "mismatched"
|
|
assert delivery["presented_paths"] == ["/mnt/user-data/outputs/old-report.md"]
|
|
assert delivery["matched_paths"] == []
|
|
assert delivery["satisfied"] is False
|
|
assert record.status == RunStatus.error
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_fenced_worker_leaves_delivery_receipt_to_peer_recovery():
|
|
"""A stale worker must not finalize the singleton delivery receipt."""
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-lease-lost")
|
|
record.ownership_lost = True
|
|
record.abort_event.set()
|
|
record.status = RunStatus.error
|
|
event_store = MemoryRunEventStore()
|
|
thread_store = SimpleNamespace(
|
|
update_display_name=AsyncMock(),
|
|
update_status=AsyncMock(),
|
|
)
|
|
on_run_completed = AsyncMock()
|
|
agent_factory = MagicMock(side_effect=AssertionError("fenced worker started the agent"))
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(
|
|
checkpointer=None,
|
|
event_store=event_store,
|
|
thread_store=thread_store,
|
|
on_run_completed=on_run_completed,
|
|
),
|
|
agent_factory=agent_factory,
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert await _delivery_events(event_store, record.thread_id, record.run_id) == []
|
|
agent_factory.assert_not_called()
|
|
thread_store.update_display_name.assert_not_awaited()
|
|
thread_store.update_status.assert_not_awaited()
|
|
on_run_completed.assert_not_awaited()
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_event_is_singleton_across_goal_continuations(monkeypatch):
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
stream_calls = 0
|
|
continuation_calls = 0
|
|
|
|
class ContinuingAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
nonlocal stream_calls
|
|
stream_calls += 1
|
|
journal = config["context"]["__run_journal"]
|
|
tool_call_id = f"call_{stream_calls}"
|
|
journal._remember_current_run_tool_calls(
|
|
AIMessage(content="", tool_calls=[{"id": tool_call_id, "name": "present_files", "args": {}}]),
|
|
caller="lead_agent",
|
|
)
|
|
artifacts = ["/mnt/user-data/outputs/report.md"]
|
|
if stream_calls == 2:
|
|
artifacts.append("/mnt/user-data/outputs/appendix.md")
|
|
journal.on_tool_end(
|
|
Command(
|
|
update={
|
|
"artifacts": artifacts,
|
|
"messages": [ToolMessage("Successfully presented files", tool_call_id=tool_call_id)],
|
|
}
|
|
),
|
|
run_id=uuid4(),
|
|
)
|
|
yield {"messages": []}
|
|
|
|
async def prepare_continuation(**kwargs):
|
|
nonlocal continuation_calls
|
|
continuation_calls += 1
|
|
if continuation_calls == 1:
|
|
return {"messages": []}
|
|
return None
|
|
|
|
monkeypatch.setattr("deerflow.runtime.runs.worker._prepare_goal_continuation_input", prepare_continuation)
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: ContinuingAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert stream_calls == 2
|
|
assert len(delivery) == 1
|
|
assert delivery[0]["content"] == {
|
|
"presented": 2,
|
|
"paths": [
|
|
"/mnt/user-data/outputs/report.md",
|
|
"/mnt/user-data/outputs/appendix.md",
|
|
],
|
|
"by_tool": {
|
|
"present_files": [
|
|
"/mnt/user-data/outputs/report.md",
|
|
"/mnt/user-data/outputs/appendix.md",
|
|
]
|
|
},
|
|
}
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_event_emitted_exactly_once_on_error_path():
|
|
run_manager = RunManager()
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
|
|
class FailingAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
raise RuntimeError("boom")
|
|
yield # pragma: no cover - make this an async generator
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=lambda *, config: FailingAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
await asyncio.sleep(0)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert len(delivery) == 1
|
|
assert delivery[0]["content"]["presented"] == 0
|
|
fetched = await run_manager.get(record.run_id)
|
|
assert fetched.status == RunStatus.error
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_is_durable_before_terminal_run_status():
|
|
events = MemoryRunEventStore()
|
|
|
|
class OrderingRunStore(MemoryRunStore):
|
|
async def update_status(self, run_id, status, *, error=None, stop_reason=None):
|
|
if status not in {"pending", "running"}:
|
|
receipt = await events.list_events("thread-1", run_id, event_types=["run.delivery"])
|
|
assert len(receipt) == 1
|
|
return await super().update_status(run_id, status, error=error, stop_reason=stop_reason)
|
|
|
|
run_store = OrderingRunStore()
|
|
run_manager = RunManager(store=run_store)
|
|
record = await run_manager.create("thread-1")
|
|
|
|
class DummyAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=events),
|
|
agent_factory=lambda *, config: DummyAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert (await run_store.get(record.run_id))["status"] == "success"
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_write_retries_before_persisting_success():
|
|
class FlakyReceiptStore(MemoryRunEventStore):
|
|
def __init__(self):
|
|
super().__init__()
|
|
self.attempts = 0
|
|
|
|
async def put_if_absent(self, **kwargs):
|
|
self.attempts += 1
|
|
if self.attempts == 1:
|
|
raise RuntimeError("transient event store outage")
|
|
return await super().put_if_absent(**kwargs)
|
|
|
|
event_store = FlakyReceiptStore()
|
|
run_store = MemoryRunStore()
|
|
run_manager = RunManager(store=run_store)
|
|
record = await run_manager.create("thread-1")
|
|
|
|
class DummyAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=event_store),
|
|
agent_factory=lambda *, config: DummyAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert event_store.attempts == 2
|
|
assert len(await _delivery_events(event_store, "thread-1", record.run_id)) == 1
|
|
assert (await run_store.get(record.run_id))["status"] == "success"
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_write_failure_preserves_real_durable_terminal_status():
|
|
class FailingReceiptStore(MemoryRunEventStore):
|
|
def __init__(self):
|
|
super().__init__()
|
|
self.attempts = 0
|
|
|
|
async def put_if_absent(self, **kwargs):
|
|
self.attempts += 1
|
|
raise RuntimeError("event store unavailable")
|
|
|
|
run_store = MemoryRunStore()
|
|
run_manager = RunManager(store=run_store)
|
|
record = await run_manager.create("thread-1")
|
|
event_store = FailingReceiptStore()
|
|
|
|
class DummyAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=event_store),
|
|
agent_factory=lambda *, config: DummyAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
# A receipt outage must not let lease recovery rewrite a genuine success
|
|
# as an error. After bounded retries, preserve the worker's real outcome.
|
|
assert event_store.attempts > 1
|
|
assert record.status == RunStatus.success
|
|
assert (await run_store.get(record.run_id))["status"] == "success"
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_produced_artifact_delivery_fails_closed_when_receipt_cannot_be_persisted(monkeypatch):
|
|
class FailingReceiptStore(MemoryRunEventStore):
|
|
async def put_if_absent(self, **kwargs):
|
|
raise RuntimeError("event store unavailable")
|
|
|
|
run_store = MemoryRunStore()
|
|
run_manager = RunManager(store=run_store)
|
|
record = await run_manager.create("thread-1")
|
|
monkeypatch.setattr(
|
|
"deerflow.runtime.runs.worker._produced_output_paths",
|
|
AsyncMock(return_value=["/mnt/user-data/outputs/report.md"]),
|
|
)
|
|
|
|
class PresentingAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
journal = config["context"]["__run_journal"]
|
|
journal._remember_current_run_tool_calls(
|
|
AIMessage(content="", tool_calls=[{"id": "call_1", "name": "present_files", "args": {}}]),
|
|
caller="lead_agent",
|
|
)
|
|
journal.on_tool_end(
|
|
Command(
|
|
update={
|
|
"artifacts": ["/mnt/user-data/outputs/report.md"],
|
|
"messages": [ToolMessage("Successfully presented files", tool_call_id="call_1")],
|
|
}
|
|
),
|
|
run_id=uuid4(),
|
|
)
|
|
yield {"messages": []}
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=FailingReceiptStore()),
|
|
agent_factory=lambda *, config: PresentingAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert record.status == RunStatus.error
|
|
assert record.error == "Artifact delivery verification failed: terminal delivery receipt could not be persisted"
|
|
assert (await run_store.get(record.run_id))["status"] == "error"
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_event_emitted_when_checkpoint_preflight_fails(monkeypatch):
|
|
run_manager = RunManager()
|
|
run_manager.update_run_completion = AsyncMock(wraps=run_manager.update_run_completion)
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
compatibility_check = AsyncMock(side_effect=RuntimeError("incompatible checkpoint"))
|
|
monkeypatch.setattr("deerflow.runtime.runs.worker.aensure_checkpoint_mode_compatible", compatibility_check)
|
|
|
|
def unexpected_agent_factory(**kwargs):
|
|
raise AssertionError("agent must not be built after preflight failure")
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=object(), event_store=store),
|
|
agent_factory=unexpected_agent_factory,
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert len(delivery) == 1
|
|
assert delivery[0]["content"] == {"presented": 0, "paths": [], "by_tool": {}}
|
|
fetched = await run_manager.get(record.run_id)
|
|
assert fetched.status == RunStatus.error
|
|
run_manager.update_run_completion.assert_not_awaited()
|
|
|
|
|
|
@pytest.mark.anyio
|
|
async def test_delivery_event_emitted_when_cancelled_waiting_for_prior_finalization(monkeypatch):
|
|
run_manager = RunManager()
|
|
run_manager.update_run_completion = AsyncMock(wraps=run_manager.update_run_completion)
|
|
record = await run_manager.create("thread-1")
|
|
store = MemoryRunEventStore()
|
|
monkeypatch.setattr(
|
|
run_manager,
|
|
"wait_for_prior_finalizing",
|
|
AsyncMock(side_effect=asyncio.CancelledError()),
|
|
)
|
|
|
|
def unexpected_agent_factory(**kwargs):
|
|
raise AssertionError("agent must not be built after preflight cancellation")
|
|
|
|
await run_agent(
|
|
_make_bridge(),
|
|
run_manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=None, event_store=store),
|
|
agent_factory=unexpected_agent_factory,
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
delivery = await _delivery_events(store, "thread-1", record.run_id)
|
|
assert len(delivery) == 1
|
|
assert delivery[0]["content"] == {"presented": 0, "paths": [], "by_tool": {}}
|
|
fetched = await run_manager.get(record.run_id)
|
|
assert fetched.status == RunStatus.interrupted
|
|
run_manager.update_run_completion.assert_not_awaited()
|