mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-09-14 16:08:41 +00:00
* 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.
439 lines
14 KiB
Python
439 lines
14 KiB
Python
"""Task lifecycle extension notifications and outcome classification."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
from types import SimpleNamespace
|
|
from unittest.mock import AsyncMock
|
|
|
|
import pytest
|
|
from deerflow_extension_api import (
|
|
EXTENSION_TASK_STORE_KEY,
|
|
ExtensionData,
|
|
TaskInfo,
|
|
TaskOutcome,
|
|
)
|
|
from langgraph.checkpoint.memory import InMemorySaver
|
|
|
|
from deerflow.extensions.notify import (
|
|
lead_task_id,
|
|
lead_task_outcome,
|
|
notify_task_start,
|
|
notify_task_stop,
|
|
subagent_task_outcome,
|
|
)
|
|
from deerflow.extensions.registry import ExtensionRegistry
|
|
from deerflow.runtime.runs.manager import RunManager
|
|
from deerflow.runtime.runs.schemas import RunStatus
|
|
from deerflow.runtime.runs.worker import RunContext, run_agent
|
|
|
|
|
|
class _Recorder:
|
|
def __init__(self) -> None:
|
|
self.events: list[tuple[str, str, str]] = []
|
|
|
|
async def on_task_start(self, app_store, task_store, info):
|
|
self.events.append(("start", info.task_id, info.kind))
|
|
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
self.events.append(("stop", info.task_id, outcome.value))
|
|
|
|
|
|
def _extensions(*contributors):
|
|
registry = ExtensionRegistry()
|
|
for index, contributor in enumerate(contributors):
|
|
with registry.attributed_to(f"ext{index}:install"):
|
|
registry.task_lifecycle(contributor)
|
|
return registry.build()
|
|
|
|
|
|
def _info(task_id: str = "task-1", kind: str = "lead") -> TaskInfo:
|
|
return TaskInfo(task_id=task_id, run_id="run-1", thread_id="thread-1", kind=kind)
|
|
|
|
|
|
def test_task_identity_and_outcome_classification_are_explicit():
|
|
assert lead_task_id("run-abc") == "run-abc"
|
|
assert lead_task_outcome(aborted=True, succeeded=True) is TaskOutcome.ABORTED
|
|
assert lead_task_outcome(aborted=False, succeeded=True) is TaskOutcome.COMPLETED
|
|
assert lead_task_outcome(aborted=False, succeeded=False) is TaskOutcome.FAILED
|
|
assert subagent_task_outcome(cancelled=True, succeeded=True) is TaskOutcome.ABORTED
|
|
assert subagent_task_outcome(cancelled=False, succeeded=True) is TaskOutcome.COMPLETED
|
|
assert subagent_task_outcome(cancelled=False, succeeded=False) is TaskOutcome.FAILED
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_start_and_stop_reach_contributors_in_order():
|
|
first = _Recorder()
|
|
second = _Recorder()
|
|
extensions = _extensions(first, second)
|
|
store = ExtensionData("task-1")
|
|
|
|
await notify_task_start(extensions, store, _info())
|
|
await notify_task_stop(extensions, store, _info(), TaskOutcome.COMPLETED)
|
|
|
|
assert first.events == [("start", "task-1", "lead"), ("stop", "task-1", "completed")]
|
|
assert second.events == first.events
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_one_malformed_or_failing_contributor_does_not_stop_the_rest():
|
|
class _WrongShape:
|
|
def on_task_start(self, app_store, task_store, info):
|
|
raise RuntimeError("sync boom")
|
|
|
|
survivor = _Recorder()
|
|
await notify_task_start(
|
|
_extensions(_WrongShape(), survivor),
|
|
ExtensionData("task-1"),
|
|
_info(),
|
|
)
|
|
|
|
assert survivor.events == [("start", "task-1", "lead")]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_notification_timeout_is_one_shared_budget():
|
|
reached: list[str] = []
|
|
|
|
class _Hang:
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
reached.append("hang")
|
|
await asyncio.sleep(10)
|
|
|
|
class _Starved:
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
reached.append("starved")
|
|
|
|
loop = asyncio.get_running_loop()
|
|
started = loop.time()
|
|
await notify_task_stop(
|
|
_extensions(_Hang(), _Starved()),
|
|
ExtensionData("task-1"),
|
|
_info(),
|
|
TaskOutcome.COMPLETED,
|
|
timeout=0.02,
|
|
)
|
|
|
|
assert reached == ["hang"]
|
|
assert loop.time() - started < 1
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_budget_exhaustion_mid_hook_logs_a_warning_not_a_traceback(caplog):
|
|
# Spending the shared budget mid-hook is the same expected operational
|
|
# condition as the pre-hook skip, not a hook failure.
|
|
class _Hang:
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
await asyncio.sleep(10)
|
|
|
|
with caplog.at_level(logging.WARNING, logger="deerflow.extensions.notify"):
|
|
await notify_task_stop(
|
|
_extensions(_Hang()),
|
|
ExtensionData("task-1"),
|
|
_info(),
|
|
TaskOutcome.COMPLETED,
|
|
timeout=0.02,
|
|
)
|
|
|
|
records = [record for record in caplog.records if record.name == "deerflow.extensions.notify"]
|
|
assert [record.levelno for record in records] == [logging.WARNING]
|
|
assert "timed out" in records[0].getMessage()
|
|
assert "task-1" in records[0].getMessage()
|
|
assert records[0].exc_info is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_contributor_timeout_error_before_budget_is_a_hook_failure(caplog):
|
|
# A TimeoutError the contributor raises on its own is not budget
|
|
# exhaustion; it stays classified as a hook failure.
|
|
class _TimedOut:
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
raise TimeoutError("the contributor's own downstream call timed out")
|
|
|
|
with caplog.at_level(logging.WARNING, logger="deerflow.extensions.notify"):
|
|
await notify_task_stop(
|
|
_extensions(_TimedOut()),
|
|
ExtensionData("task-1"),
|
|
_info(),
|
|
TaskOutcome.COMPLETED,
|
|
timeout=30,
|
|
)
|
|
# And with no notification budget at all.
|
|
await notify_task_stop(
|
|
_extensions(_TimedOut()),
|
|
ExtensionData("task-1"),
|
|
_info(),
|
|
TaskOutcome.COMPLETED,
|
|
)
|
|
|
|
records = [record for record in caplog.records if record.name == "deerflow.extensions.notify"]
|
|
assert [record.levelno for record in records] == [logging.ERROR, logging.ERROR]
|
|
assert all("failed" in record.getMessage() for record in records)
|
|
|
|
|
|
class _RunRecorder(_Recorder):
|
|
def __init__(self) -> None:
|
|
super().__init__()
|
|
self.start_infos: list[TaskInfo] = []
|
|
self.start_stores: list[ExtensionData] = []
|
|
|
|
async def on_task_start(self, app_store, task_store, info):
|
|
self.start_infos.append(info)
|
|
self.start_stores.append(task_store)
|
|
await super().on_task_start(app_store, task_store, info)
|
|
|
|
|
|
class _OkAgent:
|
|
def __init__(self) -> None:
|
|
self.runtime_context = None
|
|
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
self.runtime_context = (config or {}).get("context")
|
|
yield {"messages": []}
|
|
|
|
|
|
class _BoomAgent:
|
|
async def astream(self, graph_input, config=None, stream_mode=None, subgraphs=False):
|
|
raise RuntimeError("agent exploded")
|
|
yield # pragma: no cover
|
|
|
|
|
|
def _bridge():
|
|
return SimpleNamespace(
|
|
publish=AsyncMock(),
|
|
publish_end=AsyncMock(),
|
|
cleanup=AsyncMock(),
|
|
)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_run_agent_uses_the_run_bound_snapshot_for_lifecycle_and_task_store():
|
|
recorder = _RunRecorder()
|
|
extensions = _extensions(recorder)
|
|
manager = RunManager()
|
|
record = await manager.create("thread-ext", assistant_id="custom-agent")
|
|
agent = _OkAgent()
|
|
|
|
await run_agent(
|
|
_bridge(),
|
|
manager,
|
|
record,
|
|
ctx=RunContext(checkpointer=InMemorySaver(), extensions=extensions),
|
|
agent_factory=lambda *, config: agent,
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert record.status is RunStatus.success
|
|
assert recorder.events == [
|
|
("start", record.run_id, "lead"),
|
|
("stop", record.run_id, "completed"),
|
|
]
|
|
assert recorder.start_infos[0].agent_name == "custom-agent"
|
|
assert agent.runtime_context[EXTENSION_TASK_STORE_KEY] is recorder.start_stores[0]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_run_agent_reports_failed_and_skips_runs_that_never_started():
|
|
recorder = _RunRecorder()
|
|
extensions = _extensions(recorder)
|
|
manager = RunManager()
|
|
|
|
failed = await manager.create("thread-failed")
|
|
await run_agent(
|
|
_bridge(),
|
|
manager,
|
|
failed,
|
|
ctx=RunContext(checkpointer=InMemorySaver(), extensions=extensions),
|
|
agent_factory=lambda *, config: _BoomAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
assert failed.status is RunStatus.error
|
|
assert recorder.events[-1] == ("stop", failed.run_id, "failed")
|
|
|
|
skipped = await manager.create("thread-skipped")
|
|
await manager.cancel(skipped.run_id)
|
|
before = list(recorder.events)
|
|
await run_agent(
|
|
_bridge(),
|
|
manager,
|
|
skipped,
|
|
ctx=RunContext(checkpointer=InMemorySaver(), extensions=extensions),
|
|
agent_factory=lambda *, config: _OkAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
assert recorder.events == before
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_lead_stop_runs_after_completion_hook_and_before_stream_end():
|
|
events: list[str] = []
|
|
manager = RunManager()
|
|
record = await manager.create("thread-order")
|
|
bridge = _bridge()
|
|
bridge.publish_end.side_effect = lambda run_id: events.append("stream-end")
|
|
|
|
class _OrderingRecorder(_RunRecorder):
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
assert record.finalizing is True
|
|
assert bridge.publish_end.await_count == 0
|
|
events.append("task-stop")
|
|
await super().on_task_stop(
|
|
app_store,
|
|
task_store,
|
|
info,
|
|
outcome,
|
|
)
|
|
|
|
async def _on_completed(_record):
|
|
events.append("run-completed")
|
|
|
|
class _CancelledAgent:
|
|
async def astream(
|
|
self,
|
|
graph_input,
|
|
config=None,
|
|
stream_mode=None,
|
|
subgraphs=False,
|
|
):
|
|
raise asyncio.CancelledError()
|
|
yield # pragma: no cover
|
|
|
|
await run_agent(
|
|
bridge,
|
|
manager,
|
|
record,
|
|
ctx=RunContext(
|
|
checkpointer=InMemorySaver(),
|
|
extensions=_extensions(_OrderingRecorder()),
|
|
on_run_completed=_on_completed,
|
|
),
|
|
agent_factory=lambda *, config: _CancelledAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert events == ["run-completed", "task-stop", "stream-end"]
|
|
assert record.finalizing is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_lead_stop_interrupt_is_deferred_until_final_cleanup():
|
|
class _StopInterrupt(BaseException):
|
|
pass
|
|
|
|
class _InterruptingRecorder(_RunRecorder):
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
await super().on_task_stop(
|
|
app_store,
|
|
task_store,
|
|
info,
|
|
outcome,
|
|
)
|
|
raise _StopInterrupt("shutdown")
|
|
|
|
manager = RunManager()
|
|
record = await manager.create("thread-interrupt")
|
|
bridge = _bridge()
|
|
|
|
with pytest.raises(_StopInterrupt, match="shutdown"):
|
|
await run_agent(
|
|
bridge,
|
|
manager,
|
|
record,
|
|
ctx=RunContext(
|
|
checkpointer=InMemorySaver(),
|
|
extensions=_extensions(_InterruptingRecorder()),
|
|
),
|
|
agent_factory=lambda *, config: _OkAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert record.finalizing is False
|
|
bridge.publish_end.assert_awaited_once_with(record.run_id)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_contributor_raising_cancellederror_cannot_interrupt_run_cleanup():
|
|
# Fail-open is decided by origin, not base class: a contributor that lets a
|
|
# CancelledError escape must not skip its successors, and must not reach the
|
|
# worker's deferred-interrupt path, which would end an otherwise successful
|
|
# run as cancelled.
|
|
class _Rogue:
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
raise asyncio.CancelledError()
|
|
|
|
survivor = _RunRecorder()
|
|
manager = RunManager()
|
|
record = await manager.create("thread-rogue-stop")
|
|
bridge = _bridge()
|
|
|
|
await run_agent(
|
|
bridge,
|
|
manager,
|
|
record,
|
|
ctx=RunContext(
|
|
checkpointer=InMemorySaver(),
|
|
extensions=_extensions(_Rogue(), survivor),
|
|
),
|
|
agent_factory=lambda *, config: _OkAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
|
|
assert record.status is RunStatus.success
|
|
assert survivor.events[-1] == ("stop", record.run_id, "completed")
|
|
assert record.finalizing is False
|
|
bridge.publish_end.assert_awaited_once_with(record.run_id)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_lead_stop_cancellation_is_deferred_rather_than_swallowed():
|
|
# CancelledError derives from BaseException, so the non-fatal `except
|
|
# Exception` guard around the stop notification must not absorb it. The
|
|
# stimulus has to be a genuine cancellation of the run task — a contributor
|
|
# raising CancelledError is contained as an extension failure instead.
|
|
entered = asyncio.Event()
|
|
|
|
class _SlowRecorder(_RunRecorder):
|
|
async def on_task_stop(self, app_store, task_store, info, outcome):
|
|
await super().on_task_stop(
|
|
app_store,
|
|
task_store,
|
|
info,
|
|
outcome,
|
|
)
|
|
entered.set()
|
|
await asyncio.sleep(10)
|
|
|
|
manager = RunManager()
|
|
record = await manager.create("thread-cancel-stop")
|
|
bridge = _bridge()
|
|
|
|
task = asyncio.create_task(
|
|
run_agent(
|
|
bridge,
|
|
manager,
|
|
record,
|
|
ctx=RunContext(
|
|
checkpointer=InMemorySaver(),
|
|
extensions=_extensions(_SlowRecorder()),
|
|
),
|
|
agent_factory=lambda *, config: _OkAgent(),
|
|
graph_input={},
|
|
config={},
|
|
)
|
|
)
|
|
await entered.wait()
|
|
task.cancel()
|
|
|
|
with pytest.raises(asyncio.CancelledError):
|
|
await task
|
|
|
|
assert record.finalizing is False
|
|
bridge.publish_end.assert_awaited_once_with(record.run_id)
|