mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-09-18 10:36:17 +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.
271 lines
11 KiB
Python
271 lines
11 KiB
Python
"""Middleware for automatic thread title generation."""
|
|
|
|
import logging
|
|
import re
|
|
from typing import TYPE_CHECKING, Any, NotRequired, override
|
|
|
|
from langchain.agents import AgentState
|
|
from langchain.agents.middleware import AgentMiddleware
|
|
from langgraph.config import get_config
|
|
from langgraph.constants import TAG_NOSTREAM
|
|
from langgraph.runtime import Runtime
|
|
|
|
from deerflow.agents.middlewares.dynamic_context_middleware import is_dynamic_context_reminder
|
|
from deerflow.config.title_config import get_title_config
|
|
from deerflow.models import create_chat_model
|
|
|
|
if TYPE_CHECKING:
|
|
from deerflow.config.app_config import AppConfig
|
|
from deerflow.config.title_config import TitleConfig
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class TitleMiddlewareState(AgentState):
|
|
"""Compatible with the `ThreadState` schema."""
|
|
|
|
title: NotRequired[str | None]
|
|
|
|
|
|
class TitleMiddleware(AgentMiddleware[TitleMiddlewareState]):
|
|
"""Automatically generate a title for the thread after the first user message."""
|
|
|
|
state_schema = TitleMiddlewareState
|
|
|
|
def __init__(
|
|
self,
|
|
*,
|
|
app_config: "AppConfig | None" = None,
|
|
title_config: "TitleConfig | None" = None,
|
|
extensions=None,
|
|
):
|
|
super().__init__()
|
|
self._app_config = app_config
|
|
self._title_config = title_config
|
|
if extensions is None:
|
|
from deerflow.extensions import get_agent_build_extensions
|
|
|
|
extensions = get_agent_build_extensions()
|
|
self._extensions = extensions
|
|
|
|
def _get_title_config(self):
|
|
if self._title_config is not None:
|
|
return self._title_config
|
|
if self._app_config is not None:
|
|
return self._app_config.title
|
|
return get_title_config()
|
|
|
|
def _normalize_content(self, content: object) -> str:
|
|
if isinstance(content, str):
|
|
return content
|
|
|
|
if isinstance(content, list):
|
|
parts = [self._normalize_content(item) for item in content]
|
|
return "\n".join(part for part in parts if part)
|
|
|
|
if isinstance(content, dict):
|
|
text_value = content.get("text")
|
|
if isinstance(text_value, str):
|
|
return text_value
|
|
|
|
nested_content = content.get("content")
|
|
if nested_content is not None:
|
|
return self._normalize_content(nested_content)
|
|
|
|
return ""
|
|
|
|
@staticmethod
|
|
def _message_type(message: object) -> str | None:
|
|
message_type = getattr(message, "type", None)
|
|
if message_type is None and isinstance(message, dict):
|
|
message_type = message.get("type") or message.get("role")
|
|
if message_type == "user":
|
|
return "human"
|
|
if message_type == "assistant":
|
|
return "ai"
|
|
return message_type if isinstance(message_type, str) else None
|
|
|
|
@staticmethod
|
|
def _message_content(message: object) -> object:
|
|
if isinstance(message, dict):
|
|
return message.get("content", "")
|
|
return getattr(message, "content", "")
|
|
|
|
@staticmethod
|
|
def _is_dynamic_context_reminder_message(message: object) -> bool:
|
|
if is_dynamic_context_reminder(message):
|
|
return True
|
|
if isinstance(message, dict):
|
|
additional_kwargs = message.get("additional_kwargs")
|
|
return isinstance(additional_kwargs, dict) and bool(additional_kwargs.get("dynamic_context_reminder"))
|
|
return False
|
|
|
|
@staticmethod
|
|
def _is_user_message_for_title(message: object) -> bool:
|
|
return TitleMiddleware._message_type(message) == "human" and not TitleMiddleware._is_dynamic_context_reminder_message(message)
|
|
|
|
def _get_title_user_message(self, state: TitleMiddlewareState) -> str:
|
|
messages = state.get("messages") or []
|
|
user_msg_content = next((self._message_content(m) for m in messages if self._is_user_message_for_title(m)), "")
|
|
return self._normalize_content(user_msg_content)
|
|
|
|
def _should_generate_title(self, state: TitleMiddlewareState, *, allow_partial_exchange: bool = False) -> bool:
|
|
"""Check if we should generate a title for this thread."""
|
|
config = self._get_title_config()
|
|
if not config.enabled:
|
|
return False
|
|
|
|
# Check if thread already has a title in state
|
|
if state.get("title"):
|
|
return False
|
|
|
|
# Check if this is the first turn (has at least one user message and one assistant response).
|
|
# Defensively coerce a None ``messages`` channel (possible when reading a
|
|
# partially-initialized checkpoint) into an empty list so ``len()`` is safe.
|
|
messages = state.get("messages") or []
|
|
min_messages = 1 if allow_partial_exchange else 2
|
|
if len(messages) < min_messages:
|
|
return False
|
|
|
|
# Count user and assistant messages
|
|
user_messages = [m for m in messages if self._is_user_message_for_title(m)]
|
|
assistant_messages = [m for m in messages if self._message_type(m) == "ai"]
|
|
|
|
# Normal path: title only after first complete exchange. Interrupted path
|
|
# (``allow_partial_exchange=True``) accepts a lone first-turn user message
|
|
# so a fallback title can still be persisted when the run is cancelled
|
|
# before any AI chunk reaches the checkpoint.
|
|
return len(user_messages) == 1 and (len(assistant_messages) >= 1 or allow_partial_exchange)
|
|
|
|
def _build_title_prompt(self, state: TitleMiddlewareState) -> tuple[str, str]:
|
|
"""Extract user/assistant messages and build the title prompt.
|
|
|
|
Returns (prompt_string, user_msg) so callers can use user_msg as fallback.
|
|
"""
|
|
config = self._get_title_config()
|
|
messages = state.get("messages") or []
|
|
|
|
assistant_msg_content = next((self._message_content(m) for m in messages if self._message_type(m) == "ai"), "")
|
|
|
|
user_msg = self._get_title_user_message(state)
|
|
assistant_msg = self._strip_think_tags(self._normalize_content(assistant_msg_content))
|
|
|
|
prompt = config.prompt_template.format(
|
|
max_words=config.max_words,
|
|
user_msg=user_msg[:500],
|
|
assistant_msg=assistant_msg[:500],
|
|
)
|
|
return prompt, user_msg
|
|
|
|
def _strip_think_tags(self, text: str) -> str:
|
|
"""Remove <think>...</think> blocks emitted by reasoning models (e.g. minimax, DeepSeek-R1)."""
|
|
return re.sub(r"<think>[\s\S]*?</think>", "", text, flags=re.IGNORECASE).strip()
|
|
|
|
def _parse_title(self, content: object) -> str:
|
|
"""Normalize model output into a clean title string."""
|
|
config = self._get_title_config()
|
|
title_content = self._normalize_content(content)
|
|
title_content = self._strip_think_tags(title_content)
|
|
title = title_content.strip().strip('"').strip("'")
|
|
return title[: config.max_chars] if len(title) > config.max_chars else title
|
|
|
|
def _fallback_title(self, user_msg: str) -> str:
|
|
config = self._get_title_config()
|
|
fallback_chars = min(config.max_chars, 50)
|
|
if len(user_msg) > fallback_chars:
|
|
# Reserve room for the ellipsis so this path honours ``max_chars``
|
|
# exactly as ``_parse_title`` does on the model path.
|
|
ellipsis = "..."
|
|
body = min(fallback_chars, config.max_chars - len(ellipsis))
|
|
return user_msg[:body].rstrip() + ellipsis
|
|
return user_msg if user_msg else "New Conversation"
|
|
|
|
def _get_runnable_config(self) -> dict[str, Any]:
|
|
"""Inherit the parent RunnableConfig and add middleware tag.
|
|
|
|
This ensures RunJournal identifies LLM calls from this middleware
|
|
as ``middleware:title`` instead of ``lead_agent``.
|
|
"""
|
|
try:
|
|
parent = get_config()
|
|
except Exception:
|
|
parent = {}
|
|
config = {**parent}
|
|
config["run_name"] = "title_agent"
|
|
config["tags"] = [
|
|
*(config.get("tags") or []),
|
|
"middleware:title",
|
|
TAG_NOSTREAM,
|
|
]
|
|
return config
|
|
|
|
def _generate_title_result(self, state: TitleMiddlewareState, *, allow_partial_exchange: bool = False) -> dict | None:
|
|
"""Generate a local fallback title without blocking on an LLM call."""
|
|
if not self._should_generate_title(state, allow_partial_exchange=allow_partial_exchange):
|
|
return None
|
|
|
|
user_msg = self._get_title_user_message(state)
|
|
return {"title": self._fallback_title(user_msg)}
|
|
|
|
async def _agenerate_title_result(
|
|
self,
|
|
state: TitleMiddlewareState,
|
|
*,
|
|
task_store=None,
|
|
) -> dict | None:
|
|
"""Generate a configured LLM title asynchronously and fall back locally."""
|
|
if not self._should_generate_title(state):
|
|
return None
|
|
|
|
config = self._get_title_config()
|
|
if not config.model_name:
|
|
user_msg = self._get_title_user_message(state)
|
|
return {"title": self._fallback_title(user_msg)}
|
|
|
|
user_msg = self._get_title_user_message(state)
|
|
|
|
try:
|
|
prompt, user_msg = self._build_title_prompt(state)
|
|
# attach_tracing=False because ``_get_runnable_config()`` inherits
|
|
# the graph-level RunnableConfig (set in ``_make_lead_agent``) whose
|
|
# callbacks already carry tracing handlers; binding them again at
|
|
# the model level would emit duplicate spans.
|
|
model_kwargs = {"thinking_enabled": False, "attach_tracing": False}
|
|
if self._app_config is not None:
|
|
model_kwargs["app_config"] = self._app_config
|
|
model = create_chat_model(name=config.model_name, **model_kwargs)
|
|
invoke_config = self._get_runnable_config()
|
|
|
|
from deerflow_extension_api import SystemOperationKind
|
|
|
|
from deerflow.extensions.notify import observe_system_model_call
|
|
|
|
response = await observe_system_model_call(
|
|
self._extensions,
|
|
SystemOperationKind.TITLE,
|
|
messages=prompt,
|
|
model_name=config.model_name,
|
|
invoke_config=invoke_config,
|
|
invoke=lambda: model.ainvoke(prompt, config=invoke_config),
|
|
task_store=task_store,
|
|
)
|
|
title = self._parse_title(response.content)
|
|
if title:
|
|
return {"title": title}
|
|
except Exception:
|
|
logger.debug("Failed to generate async title; falling back to local title", exc_info=True)
|
|
return {"title": self._fallback_title(user_msg)}
|
|
|
|
@override
|
|
def after_model(self, state: TitleMiddlewareState, runtime: Runtime) -> dict | None:
|
|
return self._generate_title_result(state)
|
|
|
|
@override
|
|
async def aafter_model(self, state: TitleMiddlewareState, runtime: Runtime) -> dict | None:
|
|
from deerflow_extension_api import task_store_from_runtime
|
|
|
|
return await self._agenerate_title_result(
|
|
state,
|
|
task_store=task_store_from_runtime(runtime),
|
|
)
|