Lifecycle and Observers
Middleware sees each model and tool call inside an agent. The four contributions in this chapter see the events around and beside those calls: a task starting and stopping, DeerFlow’s own model calls, an agent being assembled, and context being compacted. None of them can change what the host does. All of them fail open under the rules in Runtime Model.
| Contribution | Registry method | Called | Sync or async | Store it receives |
|---|---|---|---|---|
TaskLifecycleContributor | registry.task_lifecycle() | Start and stop of every lead run and subagent | async, awaited | The task store |
SystemModelCallObserver | registry.system_model_observer() | After each DeerFlow-owned model call | async | The task store, or detached |
AgentAssemblyObserver | registry.agent_assembly_observer() | At the end of every agent construction | sync | The app store only |
ContextCompactionObserver | registry.context_compaction_observer() | After each summarization | async, fire-and-forget | Detached |
All the examples on this page come from one extension that registers all four:
@extension(api="0.2.0", name="observers")
def install(registry: ExtensionRegistry, config: Mapping[str, Any]) -> None:
registry.task_lifecycle(TaskTimer())
registry.system_model_observer(SystemCallLogger())
registry.agent_assembly_observer(AssemblyDriftWatcher())
registry.context_compaction_observer(CompactionLogger())Task lifecycle
class TaskLifecycleContributor(Protocol):
async def on_task_start(self, app_store: ExtensionData, task_store: ExtensionData, info: TaskInfo) -> None: ...
async def on_task_stop(self, app_store: ExtensionData, task_store: ExtensionData, info: TaskInfo, outcome: TaskOutcome) -> None: ...A task is one lead run or one subagent execution. task_store is created for that task just before on_task_start and is the same object passed to its on_task_stop and to every middleware call inside it. This makes the pair the natural place to set up and fold up per-task state.
Timing
For a lead run:
- The run is admitted and marked started. A run cancelled before this point gets neither hook.
on_task_startis awaited, before the agent graph is built.- The agent runs, including any goal continuations.
- The host persists the run’s status and token usage, syncs the thread title and status, and runs its own completion hook.
on_task_stopis awaited. The run’s finalizing barrier is still held, so a follow-up run on the same thread cannot start its lifecycle until your hook returns.- The barrier is released and the stream end is published to clients.
For a subagent, on_task_start is awaited before the subagent’s first step, and on_task_stop in its cleanup path after the sandbox lease is released, whatever the outcome.
Both hooks share the 3-second notification budget described in Runtime Model. Because on_task_stop runs before the stream end, a slow stop hook delays the moment clients see the run finish.
TaskInfo
| Field | Lead run | Subagent |
|---|---|---|
task_id | The run id | The subagent execution id |
run_id | The run id | The parent run’s id |
thread_id | The thread id | The parent thread id |
kind | "lead" | "subagent" |
parent_task_id | None | The parent run id, which is the lead task’s task_id |
agent_name | The assistant or custom agent id | The subagent name, such as general-purpose |
resumed | False | False |
resumed is part of the contract but the current host never sets it to
True. Do not rely on it to detect continuations yet.
For a subagent, task_store.scope_id is the delegating tool-call id when there is one, which is not necessarily equal to info.task_id. Use info.task_id as the task’s identity.
A subagent whose executor has no run_id, which happens under a standalone LangGraph Server or direct factory calls, skips both hooks and logs a debug line.
TaskOutcome
| Outcome | Lead run | Subagent |
|---|---|---|
completed | Status success | Status completed |
aborted | The run was stopped, or its status is interrupted | Status cancelled |
failed | Anything else, such as error | Anything else, including timeouts |
The mapping is deliberately conservative. A subagent that hit its token or turn budget can still be completed.
Example
@dataclass
class RunClock:
started: float
@dataclass
class OutcomeTally:
counts: dict[str, int] = field(default_factory=dict)
_lock: Lock = field(default_factory=Lock, repr=False)
def add(self, kind: str, outcome: TaskOutcome) -> None:
with self._lock:
key = f"{kind}:{outcome.value}"
self.counts[key] = self.counts.get(key, 0) + 1
class TaskTimer:
async def on_task_start(self, app_store: ExtensionData, task_store: ExtensionData, info: TaskInfo) -> None:
import time
task_store.set(RunClock(time.monotonic()))
async def on_task_stop(
self,
app_store: ExtensionData,
task_store: ExtensionData,
info: TaskInfo,
outcome: TaskOutcome,
) -> None:
import time
clock = task_store.get(RunClock)
elapsed = time.monotonic() - clock.started if clock is not None else float("nan")
app_store.get_or_init(OutcomeTally, OutcomeTally).add(info.kind, outcome)
logger.info("%s %s in thread %s ended %s after %.2fs", info.kind, info.task_id, info.thread_id, outcome.value, elapsed)Per-task state goes into task_store and disappears with the task; the aggregate goes into app_store. Always handle a missing value in on_task_stop: if another contributor spent the budget, your on_task_start may have been skipped.
System model calls
class SystemModelCallObserver(Protocol):
async def on_system_model_call(
self,
app_store: ExtensionData,
task_store: ExtensionData,
kind: SystemOperationKind,
request: SystemModelRequest,
result: SystemModelResult,
) -> None: ...DeerFlow makes some model calls for itself, outside the agent’s model-call chain, so middleware never sees them. This observer reports them:
kind | Call | How it is reported | Store |
|---|---|---|---|
goal | Goal-completion evaluation after an agent turn | Awaited inline | The lead task store |
title | Thread title generation | Awaited inline | The current task store |
summarization | Each summary model attempt, including fallback models | Awaited inline | The current task store |
memory | Memory extraction by the memory worker | Submitted to the notification loop, not awaited | Usually detached |
Only the async path of summarization is observed. The sync compact_state path is not reachable from the Gateway runtime and reports nothing.
Payload
SystemModelRequest is taken before the call:
| Field | Meaning |
|---|---|
messages | Always a tuple. Goal and memory pass a message list; title and summarization pass one prompt string, which becomes a one-item tuple |
model_name | The model the call used, when known |
invoke_config | The call’s runnable config when it is a mapping, else None |
SystemModelResult is taken after it:
| Field | Meaning |
|---|---|
response | The provider response on success, else None |
error | The exception on failure or cancellation, else None |
duration_ms | Wall time of the call |
The messages normalization matters: without it, iterating a prompt string would walk its characters.
Terminal paths
Every terminal path is reported, and the host sees the call’s own result or exception unchanged:
- Success and failure notify the observers inline, after the call returns or raises.
- Cancellation is routine. Stopping a run, or sending a follow-up that interrupts it, cancels in-flight goal and summarization calls after the provider tokens are spent. Awaiting observers at that point would be interrupted by a repeated cancel, so the host submits the notification to the notification loop without waiting and re-raises the cancellation.
result.erroris theCancelledError. A host with no registered loop, or one that is shutting down, drops these observations.
Inline notifications for goal, title, and summarization are awaited
without a time budget, on the path of the run. A slow observer delays the
title, the summary, or the goal decision it observes. Keep this observer to
counting and logging, and hand anything slower to a service.
Example
class SystemCallLogger:
async def on_system_model_call(
self,
app_store: ExtensionData,
task_store: ExtensionData,
kind: SystemOperationKind,
request: SystemModelRequest,
result: SystemModelResult,
) -> None:
status = "failed" if result.error is not None else "ok"
logger.info(
"system %s call on %s: %s in %.0f ms (%d message(s), scope %s)",
kind.value,
request.model_name,
status,
result.duration_ms or 0.0,
len(request.messages),
task_store.scope_id,
)system title call on gpt-4o-mini: ok in 612 ms (1 message(s), scope 7f3c...)
system goal call on gpt-4o-mini: ok in 890 ms (2 message(s), scope 7f3c...)Agent assembly
class AgentAssemblyObserver(Protocol):
def on_agent_assembled(self, app_store: ExtensionData, descriptor: AgentAssemblyDescriptor) -> None: ...When the host builds an agent it decides the effective model, renders the system prompt, filters tools through authorization, and composes the middleware stack, all inside one synchronous call. None of that is recoverable afterwards. The host therefore emits an AgentAssemblyDescriptor at the end of every construction, which normally means once per lead run and once per subagent execution.
This is the only synchronous contribution: agent construction is synchronous, and there is no event loop to await on. The observer must be cheap and must not block. It receives only the app store. When no assembly observer is registered, the host skips building descriptors entirely.
Descriptor
| Field | Meaning | In fingerprint |
|---|---|---|
namespace | "deerflow" | yes |
agent_name | lead-agent, a custom agent name, bootstrap, or the subagent name | yes |
requested_model | The model the caller asked for, if any | no |
effective_model | The model that reaches the provider | yes |
model_parameters | Behavior-affecting model settings; identity and presentation fields are excluded | yes |
thinking_enabled, reasoning_effort | The resolved reasoning settings | yes |
base_prompt_hash | canonical_hash of the rendered system prompt | yes |
tools | One ToolDescriptor per bound tool: name, description_hash, schema_hash, source, mcp_server, mcp_transport | yes, sorted by name |
middlewares | One MiddlewareDescriptor per stack entry: name, module, policy_parameters, extension | yes, in stack order |
deferred_tool_names | Tools hidden behind tool search | yes, sorted |
enabled_skills | Enabled skill names | yes, sorted |
effective_policies | Limits such as recursion limit, prompt template id, and a skill-catalog hash | yes |
build | package_version, image_digest, git_commit of the host | no |
descriptor.fingerprint is a SHA-256 over the fields marked yes. It answers “did anything about how this agent behaves change?”:
- Tools and skills are sorted because their assembly order is incidental. Middleware keeps stack order because order decides what wraps what.
buildis excluded so a redeploy of an unchanged configuration keeps every fingerprint. Comparebuilddirectly when you need to know the host changed.requested_modelis excluded because onlyeffective_modelreaches the provider.- A contributed middleware is described by the class it wraps, and its
extensionfield names the contributing entry point, so two extensions’ middleware never collapse into one entry.
image_digest and git_commit come from the DEER_FLOW_IMAGE_DIGEST and DEER_FLOW_GIT_COMMIT environment variables and read unknown when unset.
Declaring your middleware’s policy
By default the host describes a middleware by probing a fixed set of public attributes. Declare the parameters that change your middleware’s behavior instead, so a change to them changes the fingerprint:
class ToolTimer(AgentMiddleware):
def __init__(self, slow_ms: float) -> None:
super().__init__()
self.slow_ms = slow_ms
def release_policy_parameters(self) -> dict[str, object]:
return {"slow_ms": self.slow_ms}The descriptor then records MiddlewareDescriptor(name="ToolTimer", ..., policy_parameters={"slow_ms": 200}, extension="deerflow_extension_hello:install"). Values must be JSON-serializable: hash long text rather than embedding it. The contract package also exports the helpers the host uses, so an extension computes identical hashes:
canonical_json(value): JSON with sorted keys and no insignificant whitespace. RaisesTypeErroron a value it cannot serialize instead of falling back torepr.canonical_hash(value): the SHA-256 hex digest ofcanonical_json(value).collect_release_policies(middlewares): every declaration in a stack, keyed by class name. A repeated class getsName#2and so on, and a declaration that raises is recorded as{"error": "<Type>"}instead of being dropped.
Example
@dataclass
class Fingerprints:
by_agent: dict[str, str] = field(default_factory=dict)
_lock: Lock = field(default_factory=Lock, repr=False)
def swap(self, agent: str, fingerprint: str) -> str | None:
with self._lock:
previous = self.by_agent.get(agent)
self.by_agent[agent] = fingerprint
return previous
class AssemblyDriftWatcher:
def on_agent_assembled(self, app_store: ExtensionData, descriptor: AgentAssemblyDescriptor) -> None:
previous = app_store.get_or_init(Fingerprints, Fingerprints).swap(descriptor.agent_name, descriptor.fingerprint)
if previous is not None and previous != descriptor.fingerprint:
logger.warning("agent %s changed: %s -> %s", descriptor.agent_name, previous[:12], descriptor.fingerprint[:12])An exception from the observer is logged and the agent is built anyway.
Context compaction
class ContextCompactionObserver(Protocol):
async def on_context_compacted(self, app_store: ExtensionData, task_store: ExtensionData, event: CompactionEvent) -> None: ...Summarization replaces many messages with one summary. Afterwards, nothing in state records which messages became that summary. The host captures that mapping at the only moment it still exists: just before the summary call it hashes each message about to be removed, and after a summary is produced it emits a CompactionEvent.
| Field | Meaning |
|---|---|
transform_kind | "summarization" |
transform_version | "1" |
source_content_hashes | canonical_hash(message.content) for each removed message, in order |
output_content_hash | canonical_hash of the summary text |
compacted_message_count | How many messages were removed |
kept_message_count | How many messages were kept |
To match an event against messages you hold, hash exactly the same way: canonical_hash(message.content), passing the content itself. Never stringify it first, because multimodal content is a list of dicts and str() depends on key order. Do not try to rediscover the summary later by hashing what the model is shown: the prompt carries a bounded, escaped rendering of the summary, whose hash will not match output_content_hash.
The notification is fire-and-forget. It is dispatched to the notification loop without blocking the model turn, and the observer receives a detached store, because there is no live task at that call site. When no compaction observer is registered, the host skips the hashing pass too.
Example
class CompactionLogger:
async def on_context_compacted(self, app_store: ExtensionData, task_store: ExtensionData, event: CompactionEvent) -> None:
logger.info(
"%s v%s folded %d message(s) into summary %s, kept %d",
event.transform_kind,
event.transform_version,
event.compacted_message_count,
event.output_content_hash[:12],
event.kept_message_count,
)Pitfalls
- Doing I/O in a hook. Lifecycle hooks share a 3-second budget; inline system-model notifications have none and sit on the run’s path; assembly observers block construction. Buffer in a store and flush from a service.
- Keeping state on a detached store. It is discarded after the notification. Use the app store for anything that must survive.
- Assuming start implies stop. A skipped start (budget spent) or a run cancelled before it started produce asymmetric sequences. Make
on_task_stoptolerate missing state. - Treating the fingerprint as a deployment id. It deliberately ignores
build. Readdescriptor.buildto identify the host binary.