跳到正文
DEERFLOW BOOK
05 / 05EN
05

第 5 章

连续性

继续执行、继续对话、重连与审计

系统说自己“有记忆”还远远不够。DeerFlow 分别用 Checkpoint 恢复执行状态,用 Summary 和 Memory 保留对话信息,用 StreamBridge 处理断线重连,再用 RunEventStore 与 RunJournal 记录一次 Run 到底发生了什么。

CONTROL FLOW中断之后,系统可能要恢复四种东西

这些机制可能使用同一个 thread_id,但保存的内容不同,谁也不能代替谁。

  1. 01继续执行

    ThreadState 与 Checkpoint 恢复图状态,并按 reducer 继续合并

  2. 02继续对话

    Summary 保留线程语义,Memory 跨线程沉淀用户事实

  3. 03断线重连

    StreamBridge 补实时尾部,gap 时重载 durable state

  4. 04审计解释

    RunEventStore 与 RunJournal 记录事件、用量和交付事实

Four continuities

先弄清丢了什么,再决定用哪种恢复机制

sales-review 的页面被关闭后,用户可能提出四个不同要求:重新打开时看见刚才处理 sales-notes.pdf 的进度;在同一 Thread 继续修改 sales-review.md;在新 Thread 里仍记得用户偏好的报告格式;发生争议时解释这次 Run 调过什么工具、交付了什么文件。它们表面上都叫“记住”,恢复对象却完全不同。

继续执行依赖 ThreadStateCheckpoint:图走到哪里、messages、todos、sandbox_id、Artifact、delegation ledger 等状态需要按 thread 恢复。继续对话依赖两条路径:Summary 把同一 Thread 的旧消息压成 summary_textMemory 则从对话中抽取可跨 Thread 使用的用户或 Agent 事实。

断线重连依赖 StreamBridge。它保存一段按 Run 编号的实时事件,使 Last-Event-ID 之后的尾部能够重放;若尾部已经被裁剪,客户端必须重载 durable Thread state,再从服务端给出的新位置继续。审计则依赖 RunEventStoreRunJournal,把工具 receipt、模型用量、middleware 变化、workspace 与 run.delivery 等事实绑定到 run_id

所以不能用一句“DeerFlow 有 Memory”概括所有恢复能力。Memory 不保存 LangGraph 下一节点,Checkpoint 不替客户端缓存 SSE 帧,StreamBridge 不负责长期知识,RunEventStore 也不会成为 Agent 下一轮推理的状态。它们共享 thread_idrun_id,只是为了把相关记录对应起来,并不表示保存的是同一种信息。

[LangGraph 的官方 persistence 文档](https://docs.langchain.com/oss/python/langgraph/persistence)也把 checkpoint 绑定到 thread;[LangChain 的 memory 概念页](https://docs.langchain.com/oss/python/concepts/memory)则区分线程内短期状态与跨会话长期记忆。DeerFlow 的特别之处,是在这组框架概念外又显式加入 Run ownership、可重连 StreamBridge 和独立事件账本。

DEERFLOW MAP · MEMORY一份历史记录解决不了所有恢复问题

执行状态、对话信息、实时事件和审计记录,各有自己的写入时机、读取入口和故障影响。

NEW RUN
DynamicContextuser_id + agent_name
Model Context隐藏用户层 Memory 文本
MemoryMiddleware原始 messages
RECALLget_context
add / add_nowaitCAPTURE
LONG-TERM BACKEND
DeerMem默认结构化事实
OpenViking可选远程 Session
Backend contractread / write failure policy
Execution state

ThreadState 不只保存数据,还规定状态如何更新

同一 Thread 的状态并不只是 messages。固定提交的 ThreadState 还包含 sandbox、thread_data、title、artifacts、todos、goal、uploaded_filesviewed_images、promoted、delegations、skill_contextsummary_textbackground_tasks。Middleware 可以贡献兼容的 state schema,因此最终字段集合属于装配结果。

不同字段的更新规则并不相同,不能一律用“新值覆盖旧值”。messages 要按消息 ID 追加、替换或执行 RemoveMessage;artifacts 要保持顺序并去重;todos 的空列表表示明确清空;catalog hash 改变时,promoted 要整体替换;delegations 一旦进入 terminal 状态,就不能被较旧的 non-terminal 更新改回去。

Sandbox 更严格:多个 lazy tool 可以在同一 graph step 中写入相同 sandbox_id,这属于幂等合并;若同一 Thread 出现两个不同 id,系统必须 fail closed。任选一个值虽然能让图继续跑,却会破坏文件路径与执行环境的身份连续性。

agents/thread_state.py · idempotent sandbox identity reducer
def merge_sandbox(existing: SandboxState | None, new: SandboxState | None) -> SandboxState | None:
    """Reducer for sandbox state - accepts idempotent writes only.

    Multiple sandbox tools can initialize lazily in the same graph step and
    emit the same sandbox_id via Command(update=...). LangGraph needs an
    explicit reducer for that shared state key. Different sandbox ids in the
    same thread indicate a lifecycle/isolation bug, so fail closed instead of
    choosing one silently.
    """
    if new is None:
        return existing
    if existing is None:
        return new

    existing_id = existing.get("sandbox_id")
    new_id = new.get("sandbox_id")
    if existing_id == new_id:
        return existing
    raise ValueError(f"Conflicting sandbox state updates: {existing_id!r} != {new_id!r}")


SandboxStateField = Annotated[NotRequired[SandboxState | None], merge_sandbox]

代码中的 existing_idnew_id 相同才接受复写;冲突直接抛错。SandboxStateField 把 merge_sandbox 注册为 channel reducer,因此这个不变量在并发写入合并处执行,而不是靠每个工具自觉遵守。

Checkpoint 保存的是 reducer 处理后的状态。恢复时不能只是“把 JSON 读回来”,还必须使用同一份 schema 解释它;否则,未知的 middleware channel 可能被丢弃,需要直接替换的字段也可能错误地再合并一次。

这也解释为什么 Subagent 复用 ThreadState 类型却仍不是 child Thread:它以 checkpointer=False 运行,没有独立 checkpoint lineage。类型契约相同,只说明字段和 reducer 兼容,不等于拥有可恢复的持久身份。

Checkpoint recovery

所有 Checkpoint 读写都经过同一个入口

DeerFlow 支持 fulldelta 两种 checkpoint channel mode。full 保存完整 channel value;delta 依赖 DeltaChannel 的分步写入和周期快照。两者不是同一数据库内容的无差别读取方式:full 进程若把 delta checkpoint 的 sentinel 当成空状态,会形成静默数据损坏。

因此模式在进程构图时冻结,并写入 checkpoint metadata。delta 可以读取旧的 full checkpoint,形成 fulldelta 的迁移方向;反方向默认 fail closed。改变 mode 或 snapshot frequency 需要重启并保证共享同一 checkpoint 数据库的进程一致。

CheckpointStateAccessor 是读取和修改状态的统一入口。它把 effective graph、checkpointer 与已经确定的 mode 放在一起,为每次 config 加上模式标记,再由 graph 还原 state。history、update 及其异步版本都走这条路径。rollback 和 context compaction 则使用只修改状态的 graph,避免一次状态替换意外启动 Agent 节点。

runtime/checkpoint_state.py · mode-aware state access
@dataclass
class CheckpointStateAccessor:
    graph: Any
    checkpointer: Any
    mode: CheckpointChannelMode

    @classmethod
    def bind(
        cls,
        graph: Any,
        checkpointer: Any,
        *,
        store: Any | None = None,
        mode: CheckpointChannelMode = "full",
    ) -> CheckpointStateAccessor:
        graph.checkpointer = checkpointer
        if store is not None:
            graph.store = store
        return cls(graph=graph, checkpointer=checkpointer, mode=mode)

    def _prepare_config(self, config: dict[str, Any]) -> dict[str, Any]:
        prepared = {
            **config,
            "configurable": dict(config.get("configurable", {})),
            "metadata": dict(config.get("metadata", {})),
        }
        inject_checkpoint_mode(prepared, self.mode)
        return prepared

    def get(self, config: dict[str, Any]) -> Any:
        prepared = self._prepare_config(config)
        snapshot = self.graph.get_state(prepared)
        raise_if_snapshot_incompatible(snapshot, self.mode)
        return snapshot

    async def aget(self, config: dict[str, Any]) -> Any:
        prepared = self._prepare_config(config)
        snapshot = await self.graph.aget_state(prepared)
        raise_if_snapshot_incompatible(snapshot, self.mode)
        return snapshot

摘录里的 raise_if_snapshot_incompatible 会在 snapshot 返回给调用者之前检查来源模式;写入前也有兼容性检查,因为错误 checkpoint 一旦写出就无法撤销。CheckpointStateAccessor 因而同时保证两件事:读到的内容与当前模式兼容,并且使用了正确的 graph 来解释它。

对 sales-review 而言,同一 Thread 的后续 Run 可以从 checkpoint 读到消息、summary_text、Artifact 和 sandbox identity,再继续 graph;取消或 edit replay 还可以回到 pre-run rollback point。普通成功 Run 不需要 finalizing barrier,只有会回滚或重写历史的路径才需要隔离后继写入。

Checkpoint 不是 RunStore。RunStore 管 pending、running、terminal、owner lease 与取消请求;Checkpoint 管图状态。开发默认的 MemoryRunStore 只是当前 Gateway 进程内字典,不具备跨 worker 的 durable ownership;多实例部署要依靠共享持久 backend 的原子 admission 与 fencing。

Thread-local compaction

Summary 先成功生成,再提交历史替换

长对话不能无限塞进模型 context。DeerFlow 的 SummarizationMiddleware 把较旧消息压缩为 summary_text,同时保留活动尾部。它不是删除历史后再碰碰运气生成摘要,而是先选择 messages_to_summarizepreserved_messages,再调用摘要模型。

摘要模型可以显式配置;失败时可回退到本次 Run 的模型。自动压缩若所有候选失败,会返回 None,让原状态保持不变并等待未来重试;手动 /compact 可以选择 raise_on_failure,把“没有需要压缩”和“生成失败”分开报告。空白摘要也被当成失败,避免用空字符串替换真实历史。

只有非空 replacement summary 已经存在,before-summarization hooks 才收到即将移除的消息,compaction observer 才记录 source/output hash。这个顺序避免摘要失败时把同一段消息反复送入 durable Memory,也避免先删后写的不可恢复窗口。

summarization_middleware.py · summarize before replacement
async def acompact_state(
    self,
    state: AgentState,
    runtime: Runtime,
    *,
    force: bool = False,
    raise_on_failure: bool = False,
) -> ContextCompactionResult | None:
    """Async counterpart of :meth:`compact_state` (see it for ``raise_on_failure``)."""
    prepared = self._prepare_compaction(state, force=force)
    if prepared is None:
        return None
    messages_to_summarize, preserved_messages, previous_summary, total_tokens = prepared
    from deerflow_extension_api import task_store_from_runtime

    source_content_hashes = self._freeze_compaction_sources(messages_to_summarize)
    summary = await self._asummarize_with(
        messages_to_summarize,
        previous_summary=previous_summary,
        task_store=task_store_from_runtime(runtime),
    )
    if summary is None:
        if raise_on_failure:
            raise SummaryGenerationError("summary generation failed")
        return None
    # Fire hooks only once a replacement summary exists (see compact_state).
    self._fire_hooks(messages_to_summarize, preserved_messages, runtime)
    self._record_compaction(
        source_content_hashes,
        summary=summary,
        compacted_message_count=len(messages_to_summarize),
        kept_message_count=len(preserved_messages),
    )
    return ContextCompactionResult(
        summary_text=summary,
        messages_to_summarize=tuple(messages_to_summarize),
        preserved_messages=tuple(preserved_messages),
        total_tokens=total_tokens,
    )

def _maybe_summarize(self, state: AgentState, runtime: Runtime) -> dict | None:
    result = self.compact_state(state, runtime, force=False)
    if result is None:
        return None
    return {
        "messages": [
            RemoveMessage(id=REMOVE_ALL_MESSAGES),
            *result.preserved_messages,
        ],
        "summary_text": result.summary_text,
    }

async def _amaybe_summarize(self, state: AgentState, runtime: Runtime) -> dict | None:
    result = await self.acompact_state(state, runtime, force=False)
    if result is None:
        return None
    return {
        "messages": [
            RemoveMessage(id=REMOVE_ALL_MESSAGES),
            *result.preserved_messages,
        ],
        "summary_text": result.summary_text,
    }

随后 _amaybe_summarize 才写入 RemoveMessage(id=REMOVE_ALL_MESSAGES)、preserved_messagessummary_text。代码先返回 ContextCompactionResult,再构造替换 update;REMOVE_ALL_MESSAGES 是提交历史替换的动作,不是摘要生成的开始。

被压缩的语义怎样回到下一次模型调用?DurableContextMiddleware 从 state 读取 summary_text、delegations 和 skill_context。它把静态 authority contract 放在 SystemMessage,把摘要、委派结果和 Skill 描述等不可信数据放在隐藏 HumanMessage 中,避免历史数据获得 system role 权限。

因此,Summary 负责压缩同一 Thread 的旧对话,DurableContextMiddleware 则在每次 model request 前把摘要重新放回上下文。这样模型仍能看见长期任务脉络,但它不等于跨 Thread 的个性化 Memory,也不能逐字恢复每一句旧话。

Cross-thread memory

Memory 开始时读取,结束后更新

Memory 处理另一种问题:用户在新 Thread 启动下一次销售分析时,系统是否还知道“报告先给结论、金额用人民币”这类可复用偏好。MemoryManager 按 agent_nameuser_id 分桶,thread_id 作为对话来源;后端可以是 DeerMem,也可以实现同一 contract 的第三方系统。

在 middleware mode 下,DynamicContextMiddleware 通过 before_agent 进入,而不是在每次 model call 前重新检索 Memory。第一次对话入口没有既有 reminder 时,它读取 memory context 与日期,使用 ID-swap 将隐藏的 system date 和 user-role memory block 固化进 ThreadState

同一天的后续 Run 仍会经过 before_agent hook,但若 ThreadState 中已经保存 reminder,就不会再次注入;跨过午夜时只补一次日期更新。这样做能稳定 prompt prefix cache。也就是说,Memory 不是每次 model call 都重新读取的变量;即使外部 Memory 后来更新,当前 Thread 也不会在每个模型回合自动刷新。

dynamic_context_middleware.py · one agent-entry snapshot
@override
def before_agent(self, state, runtime: Runtime) -> dict | None:
    result = self._inject(state, runtime)
    self._record_effective_memory(state, result, runtime)
    return result

@override
async def abefore_agent(self, state, runtime: Runtime) -> dict | None:
    # _inject() performs synchronous file I/O (memory JSON loading) and
    # potentially blocking network calls (tiktoken encoding download on
    # first use).  Offload to a thread so the event loop is never blocked
    # — a blocking call here starves all concurrent HTTP handlers (auth,
    # SSE heartbeats, etc.).  See issue #3402.
    #
    # Bounded timeout: if startup warm-up failed silently (e.g. network
    # blip during deploy), the first request's cold tiktoken download can
    # block for tens of minutes (OS TCP timeout).  Time-box injection so
    # the request degrades gracefully (no new dynamic-context update)
    # rather than hanging. Frozen context already in state remains active.
    try:
        result = await asyncio.wait_for(
            asyncio.to_thread(self._inject, state, runtime),
            timeout=_INJECT_TIMEOUT_SECONDS,
        )
    except TimeoutError:
        logger.warning(
            "DynamicContextMiddleware: injection timed out (%.1fs); skipping new memory/date injection for this turn",
            _INJECT_TIMEOUT_SECONDS,
        )
        self._record_effective_memory(state, None, runtime)
        return None
    self._record_effective_memory(state, result, runtime)
    return result

摘录中的 _INJECT_TIMEOUT_SECONDS 还给冷文件读取或 tokenizer 初始化设上限:超时后跳过新的动态上下文,而 checkpoint 中已有 frozen context 继续生效。_record_effective_memory 只把实际使用内容的 SHA-256 身份写进 RunJournal,不复制可能包含用户数据的全文。

读取 Memory 的异常默认记录后返回空字符串,让本次 Run 在没有新记忆的情况下继续;只有 MemoryManagerError 且 backend_config.failure_policy.read=fail_closed 时才重新抛出。这里的 fail-open 是召回路径的配置语义,不能套用到前面 Checkpoint mode mismatch 的 fail-closed 状态保护。

Memory data 仍以 HumanMessage 身份注入,因为它来自用户历史,可能包含不可信文本;当前日期与框架元数据才使用 SystemMessage。这个角色分离比“把记忆拼进 system prompt”多一步,却能防止旧对话中的标签被提升为系统指令。

在 tool mode 下,事实留在 memory_search 后面,由模型按 query 检索;是否保留被动写入取决于 backend capability。无论哪种模式,Memory 都不提供 graph 下一节点、Run owner 或 SSE cursor,它只负责可复用语义。

Memory failure semantics

aadd 返回时,记忆可能还没有写入存储

Memory 写入发生在 MemoryMiddlewareafter_agent / aafter_agent。它拿到完整 state 后解析 thread_id、messages、user_idtrace_id,再把原始消息交给 MemoryManager;具体 backend 决定只保留用户输入和最终 AI 回复、怎样识别修正或强化信号,以及最终存成事实还是别的结构。

这个时机是一次 Agent graph 执行结束后,不是每个 tool turn 后。它可以看到本次完整对话结果,但若进程在 hook 前硬退出,本次更新没有机会入队;若 Memory 被禁用或缺少有效 thread_id,也会明确跳过。

异步路径写着 await manager.aadd(...),容易被读成 durable commit。实际上 MemoryManager 默认 aadd 只是调用同步 add;DeerMem 的 add 完成过滤和 queue admission,真正的 LLM extraction 与 storage write 由 process-local debounce queue 稍后执行。

memory_middleware.py · async manager boundary
@override
async def aafter_agent(self, state: MemoryMiddlewareState, runtime: Runtime) -> dict | None:
    """Use the manager's async boundary on LangGraph's async execution path."""
    add_args = self._resolve_add_args(state, runtime)
    if add_args is None:
        return None
    thread_id, messages, user_id, trace_id = add_args
    manager = await asyncio.to_thread(get_memory_manager)
    await manager.aadd(
        thread_id,
        messages,
        agent_name=self._agent_name,
        user_id=user_id,
        trace_id=trace_id,
    )
    return None

所以 manager.aadd 前面的 await 只表示 Middleware 等到了后端接口返回,不代表 Memory 已经持久化。换成另一种 backend 后,aadd 可以实现真正的异步持久写入;到底完成到哪一步由具体实现决定,调用方不能自行承诺“已经落盘”。

默认 DeerMem 会捕获 queue backpressure,把普通 Memory 更新降级成 skipped update,因此不会仅因队列已满就让已经生成 sales-review.md 的 Run 失败。但 MemoryMiddleware 本身没有吞掉任意 backend exception;第三方 aadd 若抛错仍可能越出 after_agent。best-effort 是 DeerMem 这条实现路径的语义,不是所有 MemoryManager 的无条件保证。

DeerMem queue

队列满时优先保住重要记忆

DeerMem 的 MemoryUpdateQueue 是进程内列表加 threading.Timer。相同 thread、user、agent 且同类写入会在 debounce 窗口合并;graceful shutdown 可以调用有界 flush_sync,但待处理项在进程硬退出时仍会丢失。

queue_max_depth 到达上限时,新来的普通、无 signal、非 emergency 项会抛 QueueFull;DeerMem.add 捕获后记录 warning,把失败降级成 skipped update。下个对话周期会重新提供完整会话,且未入队时 watermark 不前进,因此常规更新还有再次被提取的机会。

correction、reinforcement、preference 等 signal-bearing update 总是允许进入。Summary 删除消息前触发的 add_nowait 使用 bypass_watermark emergency 路径,也必须被接纳,因为这批即将消失的源消息不一定能在下一轮重放。

deermem/core/queue.py · selective backpressure
max_depth = self._config.queue_max_depth
if max_depth > 0 and not bypass_watermark and not signals and existing is None and len(self._items) >= max_depth:
    raise QueueFull(f"memory update queue is full (depth {len(self._items)} >= {max_depth}); non-signal update for thread {thread_id} rejected")

# Merge by signal union: a signal seen on any update for this key stays.
merged_signals = signals | (existing.signals if existing is not None else frozenset())
context = ConversationContext(
    thread_id=thread_id,
    messages=messages,
    agent_name=agent_name,
    user_id=user_id,
    trace_id=trace_id,
    signals=merged_signals,
    bypass_watermark=bypass_watermark,
)
if existing is not None:
    self._items = [c for c in self._items if not (queue_key(c.thread_id, c.user_id, c.agent_name) == key and c.bypass_watermark == bypass_watermark)]
self._items.append(context)
return context

源码里的 bypass_watermark 与 signals 共同绕过 QueueFull;existing 同 key 更新则替换旧 queue item 并合并 signal 集。emergency 与 normal 项使用不同匹配键并存,避免紧急 flush 覆盖普通项尚未提取的尾部。

这是一个明确的工程取舍:Memory 用来改善后续体验,但不能阻塞主 Run;重要信号优先,普通更新允许跳过,而且这条进程内队列并不冒充持久消息队列。面试时把这点讲清楚,比只说“我们做了长期记忆”更能体现对可靠性的理解。

Reconnect continuity

断线后先补事件;补不齐就重读状态

worker 将 metadata、updates、messages、custom events、error 与 end 发布到 StreamBridge;SSE endpoint 只订阅这条按 run_id 标识的事件流。客户端断开不会取消 Run,worker 仍由 Run owner 和 lease 管理。

重连时浏览器提交 Last-Event-IDMemoryStreamBridge 的 event id 带单调 sequence,保留窗口内可以从 cursor 后继续;Redis 等跨进程实现可以提供相同接口。StreamBridge 的职责是实时传输和有限 replay,不是 Checkpoint,也不等于永久事件存储。

如果请求 cursor 早于 retained watermark,Bridge 返回 StreamGap 并停止当前订阅。静默从最早可用事件继续会让 UI 把缺失中段误认为完整历史,所以 gap 必须成为显式控制分支,而不是 normal finish。

frontend/core/api/api-client.ts · reload durable state on a replay gap
// The SDK would otherwise ignore an unknown `gap` event and report a
// normal finish. Surface a custom control event to DeerFlow's hook, reload
// durable values, then explicitly follow only events newer than the
// retained tail captured by the server.
clearReconnectRun(threadId, runId);
yield {
  event: "custom",
  data: { type: "stream_replay_gap", ...gap },
};

const durableState = await client.threads
  .getState(threadId)
  .catch((error: unknown) => {
    throw new StreamReplayGapError(gap, recoveryAttempts, error);
  });
if (durableState.values != null) {
  yield { event: "values", data: durableState.values };
}

rememberReconnectRun(threadId, runId);
stream = resume(runId, gap.latest_available_event_id);

前端收到 stream_replay_gap 后先清除 optimistic message、transient history 与 subtask state,再调用 threads.getState 读取 durableState.values;随后用 latest_available_event_id 重新 join。摘录中的 durableState 负责恢复已提交图状态,StreamBridge 只继续补它之后的实时尾部。

这个恢复仍有界:连续 gap 超过上限会抛 StreamReplayGapError,避免无限 rejoin。页面可以给出恢复警告,而后端 Run 继续运行;“UI 暂时不完整”和“Agent 被取消”是两个独立状态。

对 sales-review 而言,刷新页面后可能先从 Checkpoint 看见已经登记的 messages、tasks 与 Artifact,再接收剩余 token 或 tool event。若最终文件已经写入但 Run settlement 尚未完成,UI 也不能仅凭重连成功宣布 Run success。

Audit continuity

Checkpoint 保存结果,RunEventStore 解释过程

Checkpoint 回答“Thread 现在是什么状态”,审计还要回答“它怎样变成这样”。RunEventStorethread_idrun_id、sequence 保存事件;RunJournal 在 Agent 执行期间归集模型调用、token usage、tool receipt、middleware state change、memory context hash 与 produced Artifact。

高频信息先在 Journal 中批量缓冲,进度快照另行节流;Subagent step 也使用有界 batch,避免每个片段都争抢持久化锁。run.delivery 是 terminal fact record,包含 presented paths 与 coverage verification;它与最终 RunStatus 有顺序关系,但不是同一行数据。

worker 的 finally 会先 flush 普通 journal events,再尝试幂等写入 run.delivery receipt,最后由仍持有 owner lease 的 worker 写入 Run 终态。已经失去 ownership 的旧 worker 会跳过这些持久化收尾,避免两个执行者同时宣布不同结论。

runtime/journal.py · restore failed batches before raising
async def flush(self) -> None:
    """Force flush remaining buffer. Called in worker's finally block."""
    if self._pending_flush_tasks:
        await asyncio.gather(*tuple(self._pending_flush_tasks), return_exceptions=True)
    while self._pending_progress_task is not None and not self._pending_progress_task.done():
        if self._pending_progress_delayed:
            self._pending_progress_task.cancel()
            await asyncio.gather(self._pending_progress_task, return_exceptions=True)
            self._progress_dirty = False
            self._pending_progress_delayed = False
            break
        await asyncio.gather(self._pending_progress_task, return_exceptions=True)

    while self._buffer:
        batch = self._buffer[: self._flush_threshold]
        del self._buffer[: self._flush_threshold]
        try:
            await self._store.put_batch(batch)
        except Exception:
            self._buffer = batch + self._buffer
            raise

摘录中的 self._buffer 按 threshold 切 batch;put_batch 失败时把 batch 放回缓冲区再抛出,避免一次瞬时异常在内存中无声删除未写事件。但 worker 捕获最终 flush 异常、记录 warning 后仍继续 receipt 与 terminal 尝试,因此普通审计事件并非成功 Run 的全有或全无事务。

delivery receipt 更严格:若本次确实产生 output,而 receipt 经有限重试仍无法写入,success candidate 可以降级为 error;随后 terminal persistence 仍可能独立失败。Chapter 02 的结算顺序与本章的审计存储在这里接上:每一步是 attempt,owner fencing 约束谁能尝试,不保证外部存储永不失败。

RunEventStore 也不等于 StreamBridge。前者为刷新后调试、历史 UI 和交付证明提供持久事实;后者为在线客户端提供低延迟流。事件可以同时经过两条路径,但一个实时帧没有持久化成功,不能从“用户看见过”推导为“审计库里存在”。

Continuity matrix

遇到恢复问题,先判断丢的是哪类信息

继续执行:身份是 thread_id 与 checkpoint namespace;写入发生在 graph state transition;持久性由 checkpointer backend 决定;失败会让状态无法继续或触发 fail-closed mode mismatch。它的核心对象是 ThreadState,不是用户画像。

继续对话:Summarythread_id 内的 summary_text 保存压缩语义,Memoryagent_name/user_id bucket 保存跨 Thread 事实。Summary replacement 进入 Checkpoint;DeerMem 更新先走 process-local queue,再写其 storage。两者都可能有损,且写入时机不同。

断线重连:身份是 run_idLast-Event-IDStreamBridge 在 Run 活跃期传递和短期重放,gap 后必须回到 durable Thread values。Bridge 丢失只影响实时视图,不应直接取消后端 Run。

审计解释:身份是 thread_id + run_id + sequence;RunJournal 先缓冲,RunEventStore 接收 batch 和 singleton receipt。普通日志写失败多为降级,交付 receipt 失败在产生 output 的成功候选上会影响 terminal status。

再加一条部署检查:MemoryRunStoreMemoryStreamBridge 与 DeerMem queue 都含 process-local 成分,不能因为接口叫 Store、Bridge 或 Memory 就推断跨进程耐久。真正的多 worker 语义来自所选 shared backend 及其原子操作。

Interview answer

面试时这样解释 DeerFlow 的记忆系统

回到 sales-review:ThreadState 用 reducer 保存 messages、sandbox、Artifact、delegation 与 summary_textCheckpointStateAccessor 以冻结 mode materialize 它;长历史先成功生成 Summary,再 flush 即将移除的消息,最后提交 RemoveMessage 替换。

跨 Thread Memorybefore_agent 第一次读取,并以普通 HumanMessage 的身份写入 ThreadStateafter_agent 再通过 MemoryManager 提交更新。await aadd 只保证后端接口已经返回;默认 DeerMem 仍使用进程内 debounce queue,队列拥堵时可以跳过普通更新,但会优先保留 signal 与 Summary 前的紧急更新。

页面断线由 StreamBridge 的有限 replay 处理;出现 gap 时清 transient state、重载 durable checkpoint values、再从 retained tail 继续。审计则由 RunJournalRunEventStore 记录工具、模型、Memory hash、workspace 和 run.delivery,并服从 Run owner fencing。

  • 90 秒表达:DeerFlow 不用一个 Memory 概念承担所有恢复。Checkpoint 恢复 ThreadState 和 graph 进度;Summary 压缩同一 Thread 的旧对话,DurableContext 在调用模型前把摘要放回上下文;长期 Memory 跨 Thread 保存用户信息,但默认写入只进入 best-effort queue;StreamBridge 负责实时事件与有限重放,出现 gap 时回到 durable state;RunJournal/RunEventStore 负责按 Run 记录过程。这些机制各自保存不同内容,写入时机和故障影响也不同。
  • 边界一:Summary 是线程内压缩,Memory 是跨线程语义;两者都不是 Checkpoint 的替代品。
  • 边界二:before_agent hook 每次 Run 会经过,但 frozen Memory 不在每个 model call 重读;await aadd 也不保证落盘。
  • 边界三:StreamBridge 不是持久状态,RunEventStore 不是模型上下文;用户实时看见一个事件,不代表它已经进入审计记录。