Skip to Content
DeerFlow

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.

ContributionRegistry methodCalledSync or asyncStore it receives
TaskLifecycleContributorregistry.task_lifecycle()Start and stop of every lead run and subagentasync, awaitedThe task store
SystemModelCallObserverregistry.system_model_observer()After each DeerFlow-owned model callasyncThe task store, or detached
AgentAssemblyObserverregistry.agent_assembly_observer()At the end of every agent constructionsyncThe app store only
ContextCompactionObserverregistry.context_compaction_observer()After each summarizationasync, fire-and-forgetDetached

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:

  1. The run is admitted and marked started. A run cancelled before this point gets neither hook.
  2. on_task_start is awaited, before the agent graph is built.
  3. The agent runs, including any goal continuations.
  4. The host persists the run’s status and token usage, syncs the thread title and status, and runs its own completion hook.
  5. on_task_stop is 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.
  6. 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

FieldLead runSubagent
task_idThe run idThe subagent execution id
run_idThe run idThe parent run’s id
thread_idThe thread idThe parent thread id
kind"lead""subagent"
parent_task_idNoneThe parent run id, which is the lead task’s task_id
agent_nameThe assistant or custom agent idThe subagent name, such as general-purpose
resumedFalseFalse

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

OutcomeLead runSubagent
completedStatus successStatus completed
abortedThe run was stopped, or its status is interruptedStatus cancelled
failedAnything else, such as errorAnything 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:

kindCallHow it is reportedStore
goalGoal-completion evaluation after an agent turnAwaited inlineThe lead task store
titleThread title generationAwaited inlineThe current task store
summarizationEach summary model attempt, including fallback modelsAwaited inlineThe current task store
memoryMemory extraction by the memory workerSubmitted to the notification loop, not awaitedUsually 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:

FieldMeaning
messagesAlways a tuple. Goal and memory pass a message list; title and summarization pass one prompt string, which becomes a one-item tuple
model_nameThe model the call used, when known
invoke_configThe call’s runnable config when it is a mapping, else None

SystemModelResult is taken after it:

FieldMeaning
responseThe provider response on success, else None
errorThe exception on failure or cancellation, else None
duration_msWall 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.error is the CancelledError. 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

FieldMeaningIn fingerprint
namespace"deerflow"yes
agent_namelead-agent, a custom agent name, bootstrap, or the subagent nameyes
requested_modelThe model the caller asked for, if anyno
effective_modelThe model that reaches the provideryes
model_parametersBehavior-affecting model settings; identity and presentation fields are excludedyes
thinking_enabled, reasoning_effortThe resolved reasoning settingsyes
base_prompt_hashcanonical_hash of the rendered system promptyes
toolsOne ToolDescriptor per bound tool: name, description_hash, schema_hash, source, mcp_server, mcp_transportyes, sorted by name
middlewaresOne MiddlewareDescriptor per stack entry: name, module, policy_parameters, extensionyes, in stack order
deferred_tool_namesTools hidden behind tool searchyes, sorted
enabled_skillsEnabled skill namesyes, sorted
effective_policiesLimits such as recursion limit, prompt template id, and a skill-catalog hashyes
buildpackage_version, image_digest, git_commit of the hostno

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.
  • build is excluded so a redeploy of an unchanged configuration keeps every fingerprint. Compare build directly when you need to know the host changed.
  • requested_model is excluded because only effective_model reaches the provider.
  • A contributed middleware is described by the class it wraps, and its extension field 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. Raises TypeError on a value it cannot serialize instead of falling back to repr.
  • canonical_hash(value): the SHA-256 hex digest of canonical_json(value).
  • collect_release_policies(middlewares): every declaration in a stack, keyed by class name. A repeated class gets Name#2 and 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.

FieldMeaning
transform_kind"summarization"
transform_version"1"
source_content_hashescanonical_hash(message.content) for each removed message, in order
output_content_hashcanonical_hash of the summary text
compacted_message_countHow many messages were removed
kept_message_countHow 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_stop tolerate missing state.
  • Treating the fingerprint as a deployment id. It deliberately ignores build. Read descriptor.build to identify the host binary.