区块链 区块链技术 比特币公众号手机端

从源码理解 DeepSeek Harness-Core

liumuhui 3小时前 阅读数 1 #区块链
文章标签 DeepSeek

前言

上一篇写完了Cordis, 这一篇开始写 DSH.
DSH 工程复杂,能讲的内容太多.
只能挑几个模块中的几个问题大概讲下.

Agent-Loop

在开始讲下 DSH 前先讲下 ReAct 对理解 DSH 更有帮助.
简单来说,就是你在调用 LLM 的时候, 在给 LLM 传递了要问的问题外,还额外传递了 Tool 列表.
Tool 列表会包含每个 Tool 的名称, 简介, 调用方法.

LLM 处理了你的问题,如果决定了这一轮要调用什么工具.
在回答的时候会带着一个 Tool Call.
此时会由你本地的 Runtime 调用该 Tool.
上一步 Tool 调用返回的结果会随着下一轮 LLM 调用发送过去.
不断循环直到 LLM 不再调用 Tool.

User
↓
Agent
↓
while loop
├─ LLM
├─ 思考 / 产生 Tool Call
├─ Tool
├─ Tool Result
└─ ...
↓
Final Answer

DSH 的 agent-loop 本质还是 ReAct 模型.

既然 ReAct 是主线,从 agent-loop Package 开始看整个 Core 开始是最好的.
其他的 Package 都会在这个 Loop 被涉及到.

为什么有两层while?

假设用户发送了消息A,此时 Agent 在处理用户用户并调用 LLM.
此时用户又发送了消息B,你可以选择丢弃消息,无法发送消息,覆盖操作,追加到当前轮,保存等下一轮.

DSH选择将消息先保存到 Inbox, 并提供了方法可以指定处理类型:
暂时有三种处理方法:

  • followup 保存并等待到下一轮
  • steer 如果当前有任务在跑,则在当前轮处理,在内层 while 的下一循环进行处理,没有任务在跑就自己开下一轮.
  • inject 跟steer 差不多.不过如果之前没有任务在跑,不开一下轮.

最外面的 While 用于从 Inbox 获取消息,并决定如何处理.

用户连发 A、B:
  followup(A) → kick 开始
    turn 1 处理 A
    turn() 发现 next-turn 还有 B → return true
    turn 2 处理 B
    inbox 空 → return false → kick 结束 → idle

看下代码里怎么写的
三种处理方式不同的地方只有后两个参数不同.
后续 turn 当作最外面 while, step 当作 最内层 while.
第三个参数为是否唤醒开启下一 turn.

// packages/core/agent-loop/src/agent.ts
followup(input: UserMessage): void { this.send(input, 'next-turn', true)}
steer(input: UserMessage): void {this.send(input, 'next-step', true)}
inject(input: UserMessage): void {this.send(input, 'next-step', false)}

send(message: UserMessage, target: InboxTarget, wakeup: boolean): void {
    const wakingAfterAbort = wakeup && this.phase.kind !== 'idle' && this.phase.abort.signal.aborted
    const resolvedTarget = wakingAfterAbort ? 'next-turn' : target
    this.inbox.splice(resolvedTarget, Infinity, 0, [message])
    if (wakeup) this.wakeDriver(wakingAfterAbort)
  }

一个 Step 经历了哪些?

  • 准备工作
    开始一轮 Step 前,必然要做一些准备工作.
    例如获取用户发送的消息.追加提示词.追加历史消息,追加上下文.

    这部分是由 preStep 函数负责的

private async preStep(target: InboxTarget, position: { turn: number; step: number }): Promise<PreparedStep> {
    ~~~~
    const claimed = this.inbox.claim(target, position.turn)
    const assembly = await this.loopCtx.systemPrompt.assemble(assembleContextFor(this, signal))
    signal.throwIfAborted()
    const sections = renderContextSections(assembly)
    const context = this.runtimeContext.project(joinContextSections(sections), sections)
    ~~~~
    return decision.kind === 'reject' ? decision : { ...decision, assembly }
  }
  • 保存消息
    保存每一轮的消息到Session,用于下一轮的上下文
for (const message of decision.messages) {
 this.session.append('user/message', message, { surfaceOp: 'append' })
}
  • renderPrompt 处理合并提示词

  • buildRequest 用于创建LLM请求

    • 确定本次请求的 LLM 配置
    • 给 agent/request hook 一个修改配置的机会
    • 把抽象配置绑定到具体模型适配器
    • 记录 request/header + request/context 到 Session
    • 最终组装真正给 LLM 的 request
  • 发到请求并用for循环接受返回.

  • 将收到的消息存到起来.

    • 注意 chunkSeqs.push 会将消息存到 session 中.
    • 这样前端UI就可以获取session中的消息实现流式输出
const { request, preparedCall } = await this.buildRequest(
        turn, step, assembly.tools, system, this.session.deriveMessages(), signal,
      )
      const assembler = new BlockAssembler()
      const chunkSeqs: number[] = []
      try {
        const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)
        signal.throwIfAborted()
        for await (const chunk of stream) {
          signal.throwIfAborted()
          chunkSeqs.push(this.session.append('assistant/chunk', { turn, step, chunk }).seq)
          assembler.push(chunk)
        }
        signal.throwIfAborted()
 }
  • 消息接受完毕后,将最终的完整消息又保存一次到 session
  • 查看返回的消息中是否有tool-call
  • 开启下一轮

如何处理并发多个Tool Call?

大部分 agent 都是一次调用一个 Tool, DSH 支持一次并发调用多个 Tool.
这又会带来几个问题:

  • 结果顺序如何保证?
  • 单个失败了怎么办?
  • 用户取消怎么办?

ToolCall 有两种执行模式,parallelexclusive:

  • parallel 并发
  • exclusive 单独

开始前根据第一个 tool 的执行模式来判断是否要并发执行.
如果是的话,会将之后所有的都传入.

const first = planned[next]!
    const mode = ctx.tools.executionMode(first.exec).kind
    const group = mode === 'parallel' ? planned.slice(next) : [first]
    const outcome = await runGroup(
      ctx, turn, step, group, mode, signal, acceptContext,
)

假设有A(p), B(p), C(e), D(e), E(e) 5个 Tool Call.p 代表 parallel,e 代表 exclusive

正常流程:

上面会将第一个parallel之后的 group 都传入.
调用fillPool的时候会再循环过滤一遍.如果碰到exclusive会截断.
后面还会不断的while循环,所以上面会是这样的处理流程:[p,p][e],[e][p]

fillPool 在执行过程中会while循环调用 startCall.
里面最关键是await ctx.tools[TOOL_RUNTIME_SCHEDULER].prepare(call.exec).
会返回三种结果,三种结果都会将执行结果保存到 slots:

  • dispatch

    • 正常异步情况.这里额外将 promise(实际执行tool) 保存到inFlight.
    • 这里并没有执行保存的 promise ,只是保存.
  • post-result

    • 调用 prepare 的时候, tool 已经出结果不需要再执行Promise.
    • needsPosttrue, 后面会跑 postExecute
  • final-result

    • 调用 prepare 的时候, tool 已经出结果不需要再执行Promise.
    • needsPostfalse, 后面不跑 postExecute

先调用一遍 commitReadyneedsPosttrue 的先执行一遍.
是因为后面的 commitReady 是在 while (inFlight.size > 0) 中.
如果这里不执行一遍就会导致前面结果为 post-resultfinal-result 的无法执行.

使用 Promise.race 来并发 inFlight.
Promise.race 有一个先完成就继续往下走.并删除 inFlight 中的位置.再继续补充新的.
这操作是在while循环中的.所以会不断的完成 inFlight 里面的 promise 直到所有的完成.
Promise 执行成功后会设置 slots.
commitReady 会判断,如果 slots 数组所有的数据都不为 undefined 则算完成.

总结下就是 slots 中的索引先按 toolCall 顺序写入.
Promise.race 是因为所有的toolCall 可能执行时间长短不一,又有最大并发数.
这样可以不用等待,一个执行完成就让出位置开始下一个.
最终返回toolCall结果是按slots的顺序并将结果插入session中 .
这样就可以实现并发并且返回结果也是正确的了.

ToolCall 出错:
其实出不出错都不改变流程.
toolResult 本身就带有字段指示是否出错.
结果还是保存到 slots 中和 session 中返回给上一层处理.

用户取消:
在用户取消的时候会设置 abortedtrue.
流程还是一样,但是因为abortedtrue,会导致前面的 while循环 中断.
可能有的任务还没完成,会在后面做额外处理.
会讲还没跑的任务插入成跳过.

if (aborted) {
    for (const call of group.slice(started)) appendSkippedToolCall(session, turn, step, call.block)
    return { consumed: group.length, aborted: true, concluded }
  }

Session

Session 为什么同时维护 Event Log 和 Surface?

session 对外提供服务的入口主要是 append 函数.
append 调用的参数格式格式是 this.session.append('event', { x, y: 'y' },{}).
第一个参数是 event 类型, 第二个是数据.
这个用法很像是事件广播类型,而不是单纯的保存数据.

看下 append 实现:

  • type Event 类型
  • data 要插入的数据
  • opts 用于告诉这个 Event 怎么进入 Surface
append<T extends SessionEventType>(
    type: T,
    data: SessionEventMap[T],
    ...opts: T extends SurfaceEventType ? [opts: SurfaceIntent] : []
  ): SessionEvent<T> {
  ~~~
  }

snapshotJsonValue 是用于检测传入的 Json 格式是否合法并复制返回一份新的.
deepFreeze 对值深度冻结.确保后续的修改操作会抛出错误.
validateNextevent 进行验证并将 surface 消息类型的seq插入 surfaceManager

const dataSnapshot = snapshotJsonValue(data)
const surfaceMetadataSnapshot = snapshotJsonValue(surfaceMetadata)
const event = deepFreeze({
      type,
      seq: this.log.length,
      time: Date.now(),
      data: dataSnapshot,
      ...(surfaceMetadataSnapshot as { surfaceOp?: unknown; sourceEventSeqs?: unknown }),
    } as unknown as SessionEvent<T>)
    this.surfaceManager.validateNext(event as SessionEvent)

获取订阅事件的回调

callbacks = collectSessionCallbacks(entry.emitCtx, [entry.carrier, 'session/event', ...callbackArgs])

将事件插入 log 后并通知订阅事件的回调.

this.log.push(event as SessionEvent)
this.eventsSnapshot = undefined
if (callbacks !== undefined && entry !== undefined) {
 invokeContainedSessionObservers(entry.emitCtx, 'session/event', entry.id, callbackArgs, callbacks)
}

大概流程是这样:

append(type, data, surface)
        │
        ▼
把外部数据复制成 Session 自己拥有的快照
        │
        ▼
检查 Event 本身是否合法
        │
        ▼
检查它对 Surface 的操作是否合法
        │
        ▼
构造 seq + time + data 的 immutable Event
        │
        ▼
        ┌──────────────┐
        │    COMMIT    │
        │  log.push()  │
        └──────┬───────┘
               │
               ▼
        通知 observers

log 中会保存所有类型的日志.
Surface 中保存的是当前模型上下文所需的、可派生为 MessageEvent.
暂时只有user/message, assistant/message, tool/result 三种类型.
例如你发送给 LLM 的消息, toolCall结果这类.

讲完了 append 讲下 deriveMessages
nodes 保存的是 Surface类型 的消息在 log 中的位置 seq.
函数逻辑很简单.就是将所有的 surface 消息保存到 derived 并返回.

deriveMessages(): Message[] {
    const surface = this.surface
    const nodes = surface.nodes
    const generation = surface.replaceGeneration
    if (generation !== this.derivedGeneration) {
      this.derived = []
      this.derivedNodes = 0
      this.derivedGeneration = generation
    }
    for (const seq of nodes.slice(this.derivedNodes)) {
      const msg = this.deriveEventMessage(this.log[seq]!)
      if (msg) this.derived.push(msg)
    }
    this.derivedNodes = nodes.length
    return [...this.derived]
  }

回到这节的问题, DSHSession 为什么同时维护 Event LogSurface?

两者负责的事情不同.

Event Log 保存 Agent 运行过程中产生的完整 Event, 包括 Surface 不需要的各种事件. 这样可以保留完整的执行历史, 之后可以基于这些日志恢复, 重建上下文, 也方便调试和追踪 Agent 的完整执行过程.

Surface 保存的是当前模型上下文所需要的 Eventseq. 通过这些 seq 可以快速找到对应的 Event, 再由 deriveMessages() 将它们转换成 Message, 最终提供给 LLM.

这样就不需要再单独维护一份 LLM 的 Message 数据, Event Log 作为唯一的真源, 而 SurfaceMessage 都是从它派生出来的.

这种设计的代价是 Event Log 会持续增长. 尤其流式输出的 assistant/chunk 也会被记录下来, DSH 运行时间越长, Log 占用的内存和持久化空间也会越大.

版权声明

本文仅代表作者观点,不代表区块链技术网立场。
本文系作者授权本站发表,未经许可,不得转载。

发表评论:

◎欢迎参与讨论,请在这里发表您的看法、交流您的观点。

热门