mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-08-01 10:56:02 +00:00
* feat(subagents): persist and display subagent step history (#3779) Capture both assistant turns and tool outputs during subagent execution, stream them in task_running events, and persist them as subagent.* run events so the subtask card's step timeline survives a reload. Backend: - step_events.py: pure layer (capture_step_message, build_subagent_step, subagent_run_event) shared by streaming and persistence - executor.py: capture ToolMessage outputs, not just AIMessage turns - worker.py: persist task_* custom events to RunEventStore (category "subagent" keeps them out of the thread feed; list_events backfills) Frontend: - core/tasks/steps.ts + api.ts: SubtaskStep model, messageToStep, eventsToSteps, mergeSteps, fetchSubtaskSteps - subtask card accumulates live steps and backfills on expand - carry run_id onto history content messages for the events endpoint * fix(subagents): show AI turns in subtask card + paginate step backfill (#3779) Two follow-ups to the subagent step-history feature: Problem 1 — reload backfill could silently truncate the step timeline because list_events capped at 500 events (seq-ASC) across the whole run. Add task_id filtering + an after_seq forward cursor to list_events (all three stores + abstract base + the /events route), and make fetchSubtaskSteps page through one task's subagent.step events until a short page. No schema migration: the DB filter rides the existing run-scoped index via event_metadata["task_id"]. Problem 2 — the card only rendered tool steps, so persisted AI turns were never shown. Replace toolStepsForDisplay with stepsForDisplay: interleave AI reasoning turns (with text) and tool steps by message_index, drop blank-text AI turns, and drop the trailing final-answer AI turn when completed (already shown as result). Card renders AI steps as muted clamped markdown with a sparkles icon. Tests: store task_id/after_seq filtering + pagination across memory/db/jsonl, the /events route forwarding, stepsForDisplay rules, and fetchSubtaskSteps pagination. Docs updated in both AGENTS.md. * make format * fix(subagents): capture full multi-tool step tail, batch step persistence, cap tool-call args (#3779) Address PR review findings on the subagent step-history feature: 1. executor.py streamed on stream_mode="values" and captured only messages[-1] per chunk, so a multi-tool-call turn (ToolNode appends one ToolMessage per call in a single super-step) lost all but the last tool output in both the live task_running stream and the persisted history. Replace with capture_new_step_messages, which walks the newly-appended tail (and still re-checks the trailing message on no-growth chunks so id-less in-place replacements survive). 2. worker.py persisted each step with the store's low-frequency put() (a per-thread advisory lock per call); a deep subagent (max_turns=150) emits hundreds of steps on the hot stream loop. Replace with _SubagentEventBuffer, which batches via put_batch (flush on terminal subagent.end, at FLUSH_THRESHOLD, and in the worker finally). 3. build_subagent_step capped only text; tool_calls[].args were copied verbatim, so a large write_file/bash payload produced an unbounded subagent.step row. Cap each call's serialized args at SUBAGENT_STEP_MAX_CHARS, flagged args_truncated. Tests updated/added for all three; AGENTS.md refreshed. * fix(subagents): merge backfill into latest subtask state; reuse message_content_to_text (#3779) Address the remaining two PR review findings: 4. subtask-card's fetchSubtaskSteps().then(updateSubtask) closed over a stale tasks snapshot: a late-resolving backfill wrote setTasks({...stale}), clobbering SSE steps/status and sibling subtasks that arrived during the fetch. useUpdateSubtask now reads/writes through a tasksRef mirroring the latest state (ref-to-latest), and the pure per-subtask transition is extracted to core/tasks/subtask-update.ts::computeNextSubtask (unit-tested). 5. step_events._content_to_text duplicated deerflow.utils.messages. message_content_to_text; call the shared helper instead (guarding None content with 'or ""' so a tool-call-only turn still renders as ""). Tests added for computeNextSubtask and the None-content case; AGENTS.md docs updated.
116 lines
4.0 KiB
Python
116 lines
4.0 KiB
Python
"""Abstract interface for run event storage.
|
|
|
|
RunEventStore is the unified storage interface for run event streams.
|
|
Messages (frontend display) and execution traces (debugging/audit) go
|
|
through the same interface, distinguished by the ``category`` field.
|
|
|
|
Implementations:
|
|
- MemoryRunEventStore: in-memory dict (development, tests)
|
|
- Future: DB-backed store (SQLAlchemy ORM), JSONL file store
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import abc
|
|
|
|
|
|
class RunEventStore(abc.ABC):
|
|
"""Run event stream storage interface.
|
|
|
|
All implementations must guarantee:
|
|
1. put() events are retrievable in subsequent queries
|
|
2. seq is strictly increasing within the same thread
|
|
3. list_messages() only returns category="message" events
|
|
4. list_events() returns all events for the specified run
|
|
5. Returned dicts match the RunEvent field structure
|
|
"""
|
|
|
|
@abc.abstractmethod
|
|
async def put(
|
|
self,
|
|
*,
|
|
thread_id: str,
|
|
run_id: str,
|
|
event_type: str,
|
|
category: str,
|
|
content: str | dict = "",
|
|
metadata: dict | None = None,
|
|
created_at: str | None = None,
|
|
) -> dict:
|
|
"""Write an event, auto-assign seq, return the complete record."""
|
|
|
|
@abc.abstractmethod
|
|
async def put_batch(self, events: list[dict]) -> list[dict]:
|
|
"""Batch-write events. Used by RunJournal flush buffer.
|
|
|
|
Each dict's keys match put()'s keyword arguments.
|
|
Returns complete records with seq assigned.
|
|
"""
|
|
|
|
@abc.abstractmethod
|
|
async def list_messages(
|
|
self,
|
|
thread_id: str,
|
|
*,
|
|
limit: int = 50,
|
|
before_seq: int | None = None,
|
|
after_seq: int | None = None,
|
|
) -> list[dict]:
|
|
"""Return displayable messages (category=message) for a thread, ordered by seq ascending.
|
|
|
|
Supports bidirectional cursor pagination:
|
|
- before_seq: return the last ``limit`` records with seq < before_seq (ascending)
|
|
- after_seq: return the first ``limit`` records with seq > after_seq (ascending)
|
|
- neither: return the latest ``limit`` records (ascending)
|
|
"""
|
|
|
|
@abc.abstractmethod
|
|
async def list_events(
|
|
self,
|
|
thread_id: str,
|
|
run_id: str,
|
|
*,
|
|
event_types: list[str] | None = None,
|
|
task_id: str | None = None,
|
|
limit: int = 500,
|
|
after_seq: int | None = None,
|
|
) -> list[dict]:
|
|
"""Return the full event stream for a run, ordered by seq ascending.
|
|
|
|
Optionally filter by ``event_types`` and/or ``task_id`` (matched against
|
|
``metadata["task_id"]``). ``after_seq`` is a forward cursor returning the
|
|
first ``limit`` records with seq > after_seq, so callers can page through
|
|
a single subagent task's events without the run-wide ``limit`` truncating
|
|
the tail (#3779).
|
|
"""
|
|
|
|
@abc.abstractmethod
|
|
async def list_messages_by_run(
|
|
self,
|
|
thread_id: str,
|
|
run_id: str,
|
|
*,
|
|
limit: int = 50,
|
|
before_seq: int | None = None,
|
|
after_seq: int | None = None,
|
|
) -> list[dict]:
|
|
"""Return displayable messages (category=message) for a specific run, ordered by seq ascending.
|
|
|
|
Supports bidirectional cursor pagination:
|
|
- after_seq: return the first ``limit`` records with seq > after_seq (ascending)
|
|
- before_seq: return the last ``limit`` records with seq < before_seq (ascending)
|
|
- neither: return the latest ``limit`` records (ascending)
|
|
"""
|
|
|
|
@abc.abstractmethod
|
|
async def count_messages(self, thread_id: str) -> int:
|
|
"""Count displayable messages (category=message) in a thread."""
|
|
|
|
@abc.abstractmethod
|
|
async def delete_by_thread(self, thread_id: str) -> int:
|
|
"""Delete all events for a thread. Return the number of deleted events."""
|
|
|
|
@abc.abstractmethod
|
|
async def delete_by_run(self, thread_id: str, run_id: str) -> int:
|
|
"""Delete all events for a specific run. Return the number of deleted events."""
|