deer-flow/backend/tests/test_extension_subagent_lifecycle.py
Nan Gao 7389331e65
feat(extensions): observe task lifecycle and system model calls (#4684)
* feat(extensions): observe task lifecycle and system model calls

PR 1 (#4636) gave extensions a middleware chain, and a middleware only sees
what passes through the agent graph. Two runtime surfaces stay invisible to
it: when a lead run or a subagent begins and ends, and the DeerFlow-owned
model calls made outside the graph. This slice adds both, with no new
Gateway surface -- routers, services, and the reference extension stay in
PR 3.

Contract (deerflow-extension-api 0.1.1)
---------------------------------------
Two contribution kinds join `middlewares` on the registry:
`task_lifecycle` (`on_task_start` / `on_task_stop`, receiving a `TaskInfo`
and a conservative `TaskOutcome` of completed / aborted / failed) and
`system_model_observer` (`on_system_model_call`, receiving a
`SystemOperationKind`, a `SystemModelRequest` snapshot, and a
`SystemModelResult` carrying either the response or the provider exception
plus a duration).

`SystemModelRequest.messages` normalizes to a tuple at construction. Goal
evaluation and memory extraction pass a message list while title generation
and summarization pass one prompt string, and a bare `str` already satisfies
`Sequence` -- without normalization an observer iterating `request.messages`
would silently walk characters. Copying also makes the frozen snapshot
immutable in fact rather than only by declaration, since observations may run
after the call site returns and keeps mutating its own list.

Registry marks and rollbacks become per-bucket and positional, so an
`install()` that fails after registering two different kinds cannot leave one
of them behind. `needs_task_store` now covers all three kinds: a deployment
that registers only lifecycle hooks still gets a task store.

Task lifecycle
--------------
The lead worker notifies start after the run has started and stop after
completion persistence and the completion hook, but before clearing the
finalizing barrier and publishing the stream end -- holding the barrier
across stop is what keeps a same-thread replacement run from overlapping this
task's lifecycle. Cancellation raised out of the stop notification is
deferred, not propagated in place, so a cancelled run still clears the
barrier and emits its end frame. A subagent with a parent `run_id` wraps its
execution in the same pair inside `finally`, reporting `parent_task_id` so a
delegation tree is reconstructable; a subagent without a `run_id` (embedded
client, standalone LangGraph Server) logs and skips rather than inventing a
parent. Contributors run in registration order inside one shared 3s budget
and every failure is logged and failed open.

System model calls
------------------
Four kinds cover the model calls the middleware chain cannot see: goal
evaluation, memory extraction, title generation, and summarization. Each site
reports both terminal paths without changing the provider exception the host
observes, short-circuits on `has_system_model_observers`, and passes the live
task store when the runtime has one (detached work gets an isolated store).
The sync summarization half stays unobserved on purpose -- it and its only
host caller are the sync side of an async-only runtime, so notifying there
would block a thread on a call site the host never reaches; the reason is
recorded at the call site.

The DeerMem backend must stay vendorable and cannot import the extension API,
so it reports through a new `MemoryCallbacks.on_memory_llm_result` host hook
that the DeerFlow-side callbacks translate into an observation.

Notification loop
-----------------
Extension resources must be touched on the loop that created them, but
subagents can execute on isolated loops and DeerMem runs on a worker thread.
The Gateway registers its serving loop before any runtime dependency starts
and resets it last through the exit stack, so every startup-failure and
cancellation path is covered. Awaited hooks raised on another loop are
dispatched across with `run_coroutine_threadsafe` and awaited under the same
budget; synchronous sites submit fire-and-forget work. Shutdown stops
accepting detached observations before the memory flush -- that flush runs on
a worker thread and can emit memory observations -- while keeping the loop
alive for awaited task hooks until run and subagent drain completes.

Tests
-----
`test_extension_task_lifecycle.py`, `test_extension_subagent_lifecycle.py`,
and `test_extension_system_model_calls.py` cover ordering, fail-open, budget
exhaustion, snapshot binding under a concurrent singleton replacement, the
loop-dispatch and shutdown-suspension paths, and both terminal paths at every
call site. `test_gateway_run_drain_shutdown.py` pins the stop-before-barrier
and drain ordering.

* fix(extensions): decide notification fail-open by origin, observe cancellation

`_notify_each` only guarded `Exception`, so a contributor letting a
`CancelledError` escape — an extension implementing an internal timeout with
cancellation, say — skipped its successors and reached the worker's
deferred-interrupt path, ending an otherwise successful run as cancelled.
Fail-open is about where a failure came from, not its base class: only a
genuine cancellation of the host task increments `Task.cancelling()`, so
propagate on that and contain everything else. `KeyboardInterrupt` /
`SystemExit` still propagate.

`observe_system_model_call` skipped observers on cancellation for the same
base-class reason, leaving goal / title / summarization silent on a terminal
path that is routine — interrupt/rollback admission and shutdown both cancel
the run task, with the provider tokens already spent. Awaiting observers there
is unreliable (a repeated cancel interrupts that await before any of them
runs), so report through the same non-blocking submission the synchronous
memory bridge uses, then propagate the cancellation untouched.

DeerMem keeps `BaseException` around its provider call, now with the reason
recorded: that path runs on a worker thread, where cancelling the awaiting
side never interrupts the running thread, so `CancelledError` cannot arrive
at all. Its host-hook wrapper narrows to `Exception` — only the hook's own
failures are non-fatal, and an observability path must not swallow a process
teardown signal.

* fix(extensions): warn on budget exhaustion, scope observer logs by task, propagate teardown

Review response on #4684:

- The memory observation bridge caught BaseException, which would swallow
  a teardown signal raised while dispatching; it now catches Exception,
  matching the boundary the DeerMem-side call site documents and tests.
- A notification-budget timeout raised mid-hook fell into the generic
  hook-failure path and logged an asyncio-internal traceback; it now logs
  a warning like the pre-hook budget skip, while a TimeoutError a
  contributor raises on its own stays classified as a hook failure.
- System model observer logs passed the operation kind as the task id,
  so log lines said "task goal/title/..."; they now carry the task
  scope id alongside the kind.
2026-08-11 16:33:22 +08:00

260 lines
8.9 KiB
Python

"""Subagents expose the same task lifecycle contract as lead runs."""
from __future__ import annotations
import sys
from types import ModuleType, SimpleNamespace
from unittest.mock import MagicMock
import pytest
from deerflow_extension_api import EXTENSION_TASK_STORE_KEY, ExtensionData, TaskInfo, TaskOutcome
from langchain_core.messages import AIMessage
from deerflow.extensions import reset_loaded_extensions, set_loaded_extensions
from deerflow.extensions.registry import ExtensionRegistry
_MOCKED_MODULE_NAMES = (
"deerflow.agents",
"deerflow.agents.thread_state",
"deerflow.agents.middlewares",
"deerflow.agents.middlewares.thread_data_middleware",
"deerflow.sandbox",
"deerflow.sandbox.middleware",
"deerflow.sandbox.security",
"deerflow.models",
"deerflow.skills.storage",
)
@pytest.fixture
def env():
"""Import the real executor behind conftest's cycle-breaking mock."""
reset_loaded_extensions()
original_modules = {name: sys.modules.get(name) for name in _MOCKED_MODULE_NAMES}
original_executor = sys.modules.get("deerflow.subagents.executor")
subagents_pkg = sys.modules.get("deerflow.subagents")
missing = object()
original_executor_attr = getattr(subagents_pkg, "executor", missing) if subagents_pkg is not None else missing
sys.modules.pop("deerflow.subagents.executor", None)
if subagents_pkg is not None and hasattr(subagents_pkg, "executor"):
delattr(subagents_pkg, "executor")
try:
for name in _MOCKED_MODULE_NAMES:
sys.modules[name] = MagicMock()
storage_module = ModuleType("deerflow.skills.storage")
storage_module.get_or_new_skill_storage = lambda **kwargs: SimpleNamespace(load_skills=lambda *, enabled_only: [])
storage_module.get_or_new_user_skill_storage = lambda user_id, **kwargs: SimpleNamespace(load_skills=lambda *, enabled_only: [])
sys.modules["deerflow.skills.storage"] = storage_module
from deerflow.subagents.config import SubagentConfig
from deerflow.subagents.executor import (
SubagentExecutor,
SubagentResult,
SubagentStatus,
)
sys.modules["deerflow.subagents.executor"].get_app_config = lambda: SimpleNamespace(
tool_search=SimpleNamespace(enabled=False),
authorization=SimpleNamespace(enabled=False),
)
yield SimpleNamespace(
SubagentConfig=SubagentConfig,
SubagentExecutor=SubagentExecutor,
SubagentResult=SubagentResult,
SubagentStatus=SubagentStatus,
)
finally:
reset_loaded_extensions()
for name, original in original_modules.items():
if original is None:
sys.modules.pop(name, None)
else:
sys.modules[name] = original
if original_executor is None:
sys.modules.pop("deerflow.subagents.executor", None)
else:
sys.modules["deerflow.subagents.executor"] = original_executor
subagents_pkg = sys.modules.get("deerflow.subagents")
if subagents_pkg is not None:
if original_executor_attr is missing:
if hasattr(subagents_pkg, "executor"):
delattr(subagents_pkg, "executor")
else:
setattr(subagents_pkg, "executor", original_executor_attr)
class _Recorder:
def __init__(self) -> None:
self.starts: list[TaskInfo] = []
self.stops: list[tuple[TaskInfo, TaskOutcome]] = []
self.stores: list[ExtensionData] = []
async def on_task_start(self, app_store, task_store, info):
self.starts.append(info)
self.stores.append(task_store)
async def on_task_stop(self, app_store, task_store, info, outcome):
self.stops.append((info, outcome))
self.stores.append(task_store)
def _loaded(recorder):
registry = ExtensionRegistry()
with registry.attributed_to("demo:install"):
registry.task_lifecycle(recorder)
return registry.build()
def _executor(env, **overrides):
config = env.SubagentConfig(
name="researcher",
description="d",
system_prompt="p",
tools=[],
)
kwargs = {"run_id": "run-1", "thread_id": "thread-1"}
kwargs.update(overrides)
return env.SubagentExecutor(config=config, tools=[], **kwargs)
class _CompletingAgent:
def __init__(self, seen: dict | None = None) -> None:
self.seen = seen
async def astream(self, *args, **kwargs):
if self.seen is not None:
self.seen["context"] = kwargs.get("context")
yield {"messages": [AIMessage(content="done")]}
async def _noop_initial_state(self, task):
return ({}, [], None)
@pytest.mark.asyncio
async def test_subagent_success_emits_shaped_start_and_completed_stop(monkeypatch, env):
recorder = _Recorder()
set_loaded_extensions(_loaded(recorder))
executor = _executor(env)
seen: dict = {}
monkeypatch.setattr(env.SubagentExecutor, "_build_initial_state", _noop_initial_state)
monkeypatch.setattr(
env.SubagentExecutor,
"_create_agent",
lambda self, tools, **kwargs: _CompletingAgent(seen),
)
result = await executor._aexecute("do the thing")
assert result.status is env.SubagentStatus.COMPLETED
[info] = recorder.starts
assert info == TaskInfo(
task_id=result.task_id,
run_id="run-1",
thread_id="thread-1",
kind="subagent",
parent_task_id="run-1",
agent_name="researcher",
)
assert recorder.stops == [(info, TaskOutcome.COMPLETED)]
assert recorder.stores[0] is recorder.stores[1]
assert seen["context"][EXTENSION_TASK_STORE_KEY] is recorder.stores[0]
@pytest.mark.asyncio
async def test_subagent_failure_and_cancellation_map_to_distinct_outcomes(monkeypatch, env):
recorder = _Recorder()
set_loaded_extensions(_loaded(recorder))
executor = _executor(env)
async def _fail_before_agent(self, task):
raise RuntimeError("build failed")
monkeypatch.setattr(env.SubagentExecutor, "_build_initial_state", _fail_before_agent)
failed = await executor._aexecute("fail")
assert failed.status is env.SubagentStatus.FAILED
assert recorder.stops[-1][1] is TaskOutcome.FAILED
monkeypatch.setattr(env.SubagentExecutor, "_build_initial_state", _noop_initial_state)
monkeypatch.setattr(
env.SubagentExecutor,
"_create_agent",
lambda self, tools, **kwargs: _CompletingAgent(),
)
holder = env.SubagentResult(
task_id="cancel-me",
trace_id="trace",
status=env.SubagentStatus.RUNNING,
)
holder.cancel_event.set()
cancelled = await executor._aexecute("cancel", holder)
assert cancelled.status is env.SubagentStatus.CANCELLED
assert recorder.stops[-1][1] is TaskOutcome.ABORTED
assert recorder.stops[-1][0].task_id == "cancel-me"
@pytest.mark.asyncio
async def test_subagent_base_exception_still_emits_failed_stop(monkeypatch, env):
recorder = _Recorder()
set_loaded_extensions(_loaded(recorder))
executor = _executor(env)
async def _hard_stop(self, task):
raise KeyboardInterrupt("host shutdown")
monkeypatch.setattr(env.SubagentExecutor, "_build_initial_state", _hard_stop)
with pytest.raises(KeyboardInterrupt):
await executor._aexecute("stop")
assert recorder.stops[0][1] is TaskOutcome.FAILED
@pytest.mark.asyncio
async def test_subagent_without_parent_run_skips_lifecycle_but_keeps_task_store(monkeypatch, env):
recorder = _Recorder()
set_loaded_extensions(_loaded(recorder))
executor = _executor(env, run_id=None)
seen: dict = {}
monkeypatch.setattr(env.SubagentExecutor, "_build_initial_state", _noop_initial_state)
monkeypatch.setattr(
env.SubagentExecutor,
"_create_agent",
lambda self, tools, **kwargs: _CompletingAgent(seen),
)
result = await executor._aexecute("direct")
assert recorder.starts == []
assert recorder.stops == []
assert seen["context"][EXTENSION_TASK_STORE_KEY].scope_id == result.task_id
@pytest.mark.asyncio
async def test_subagent_keeps_one_snapshot_across_build_context_and_hooks(monkeypatch, env):
first = _Recorder()
second = _Recorder()
snapshot = _loaded(first)
set_loaded_extensions(snapshot)
executor = _executor(env)
seen: dict = {}
async def _switch_singleton(self, task):
set_loaded_extensions(_loaded(second))
return ({}, [], None)
def _capture_agent(self, tools, *, deferred_setup=None, extensions=None):
seen["extensions"] = extensions
return _CompletingAgent(seen)
monkeypatch.setattr(env.SubagentExecutor, "_build_initial_state", _switch_singleton)
monkeypatch.setattr(env.SubagentExecutor, "_create_agent", _capture_agent)
await executor._aexecute("snapshot")
assert seen["extensions"] is snapshot
assert len(first.starts) == len(first.stops) == 1
assert second.starts == second.stops == []
assert seen["context"][EXTENSION_TASK_STORE_KEY] is first.stores[0]