PR #108 + #109:事件存储双减负——raw 载荷下沉文件,流式文本合并落库

figuretu/eyrie · main…feat/event-rawjson-file-offload…feat/event-delta-db-coalesce(stacked)· 2026-07-02 · 自包含,读完即弃

#108 5 commits · 28 文件 · +854/−104 · 42% 测试
#109 6 commits · 9 文件 · +1258/−36 · 66% 测试
范围 几乎全在 apps/daemon

写法说明:本文按「逐跳走读」展开——每个机制都给真实代码片段(取自分支、经裁剪,青色斜体注释为解读所加,灰色斜体是源码原注释的保留或意译),每段代码标注所在文件。每条旅程结尾有一张「排查路标」:将来出问题时,症状对应去哪个文件看哪个函数。

1TL;DR

agent_events 是 daemon 最热的写路径:agent 的每个 provider 事件都进这张 SQLite 表。两个 stacked PR 分别砍掉这条路径上的两块「写了没人读 / 写得太碎」的负担:

两个 PR 都不改任何对外接口:事件流的 envelope、回放路径、断线重连、游标语义全部保持原样。#109 base 在 #108 的分支上,先合 #108。

2变更地图

两个 PR 合计 2252 行变更,分布如下。测试占大头(56%),非测试代码里几乎没有机械变更——除 schema 删列和快照重生成外,全部是设计承载代码,需要细读。

daemon/tests
1315 行 · 58%
daemon/src
910 行 · 40%
packages/db
22 行 · 1%
docs
5 行 · <1%
PR设计重心(要细读)可放心略过
#108 services/agent-event-file-writer.ts(新,286 行:串行链、rawSeq、封口、清扫);agent/session-sink.ts 的 emit 接缝;agent/service.ts 的 purge 挂点 event-codec.ts 的 −57 行是把 rawJson 参数从十几个 encoder 签名里机械删除;migration/snapshot 重生成;各测试文件里删 rawJson: null 字段
#109 agent/message-delta-coalescer.ts(新,272 行);session-sink.ts 的 emit 改道 + 边界 flush;service.ts 三处终结点排水 无——9 个文件全是设计承载代码或钉行为的测试

3架构一图流

改动前后,一条 provider 事件从 runner 到订阅者的通路对比。关键变化:raw 载荷分叉去了文件(虚线=best-effort,失败不影响主路);message.delta 在落库前多了一个内存缓冲站。

以前 · 单通道,逐 chunk 落行

Provider Runner
每 chunk 一次 emit
SessionSink
SessionSink
每事件一事务
行内含 raw_json
agent_events (raw_json)
已落库行
逐 chunk 广播
订阅者

现在 · raw 分叉进文件,delta 先合并

SessionSink
best-effort append
每原始事件一行
<runId>.raw.jsonl
message.delta
缓冲,16KB/250ms/边界 flush
Coalescer
合并行/状态事件
一事务一行(无 raw 列)
落库后按块广播
agent_events → 订阅者

注意状态类事件(approval / input / run / control / context)完全不经过缓冲站:它们的落库和状态投影(比如把审批置为 pending)在同一个事务里,必须即刻、原子地发生。合并只作用于纯展示的 message.delta

4数据形状先行

先认识四个新形状,后面旅程里不再解释。

取证文件的一行。每行一个 JSON 对象,rawSeq 是 per-run 文件内自己的序号,和数据库的事件 seq 完全解耦——因为 #109 之后 DB 行会被合并,文件行和 DB 行不再一一对应,绑 DB seq 就是给自己埋雷:

apps/daemon/src/services/agent-event-file-writer.ts真实代码(节选)
/** One lossy-readable forensic raw event line stored beside a run. */
export type RawLine = {
  rawSeq: number    // per-run 文件内序号,独立于 DB 事件 seq
  type: string      // 归一化后的事件类型
  role?: string
  itemId?: string
  occurredAt?: number
  raw: unknown      // adapter 发出的原始 provider 载荷
}

写文件的 port。三个方法覆盖整个生命周期,注入到 sink / service / runner-manager;默认实现是 no-op(窄测试和不关心文件的装配用),生产在 daemon 启动时创建真实实现:

apps/daemon/src/services/agent-event-file-writer.ts真实代码(节选)
/** Append-only raw event file port owned by daemon agent execution. */
export interface AgentEventFileWriter {
  append(sessionId: string, runId: string, line: {...}): Promise<void>
  purgeSession(sessionId: string): Promise<void>   // 整 session 目录 rm
  sweepOrphans(isLiveSession: (id: string) => Promise<boolean>): Promise<number>
}

合并器吐出的单元。合并后的 delta 始终带着它被缓冲时所属的 run——session 可能已经推进到下一个 run,文本必须落回原 run:

apps/daemon/src/agent/message-delta-coalescer.ts真实代码(节选)
export type RunScopedMessageDelta = {
  runId: string       // 接收该 chunk 时捕获的 run,钉住归属
  sessionId: string
  event: CoalescedMessageDeltaEvent   // 合并好的 message.delta,text 是拼接结果
  createdAt: number   // flush 时刻(注入的 clock 给出)
}

缓冲分组的 key。合并按 (runId, sessionId, role, itemId) 四元组分组——不同消息、不同角色(assistant / reasoning)永远不会拼进同一行;itemId 缺失时用显式的「缺失」标签,避免和任何真实 itemId 撞车:

apps/daemon/src/agent/message-delta-coalescer.ts真实代码(节选)
// NUL 分隔符保证各段无歧义;presence 标签让「没有 itemId」区别于任何 itemId 字符串
function bufferKey(runId, sessionId, role, itemId) {
  const itemKey =
    itemId === undefined ? ITEM_ID_MISSING : `${ITEM_ID_PRESENT}${KEY_SEPARATOR}${itemId}`
  return `${runId}${KEY_SEPARATOR}${sessionId}${KEY_SEPARATOR}${role}${KEY_SEPARATOR}${itemKey}`
}

schema 侧只有减法:agent_events.raw_json 列和它的 json_valid 约束从建表 SQL、drizzle 表定义、行类型里一并消失,迁移基线整体重生成(开发期单基线策略,不是新增一条迁移)。

5旅程 A:一条事件的双写落盘(#108)

跟一条带 raw 的 provider 事件走完 emit:先分叉写文件,再权威落库,最后广播。走通后你会知道文件写坏了为什么用户毫无感知,以及进程崩溃后文件如何自愈。

全景 · 4 跳
emit 接缝
session-sink.ts
串行链发号
agent-event-file-writer.ts
appendEvent 落库(不变) broadcast(不变) 取证读
readEventRaw

A.1emit 接缝:先文件后 DB,但文件永远不能拖垮 DB

改动前,emitevent.raw 序列化成 rawJson 塞进同一个 appendEvent 事务。改动后 raw 在落库之前分叉写文件——但这条分叉被 try/catch 完全隔离:文件写失败只打一条 warn 日志,DB 写和广播照常。整个 PR 最重要的不变量就是这一条:best-effort 的 raw 写永远不能丢掉一条已提交事件或它的广播

apps/daemon/src/agent/session-sink.ts真实代码(节选)
async emit(event: AgentProviderEvent, opts?: { runId?: string }): Promise<void> {
  const runId = opts?.runId ?? this.currentRunId
  if (!runId) throw new Error('SessionSink.emit called without a current run id.')
  // message.delta 分支见旅程 C;这里走非 delta 事件的直落路径

  // input.requested 携带未脱敏的请求体,可能含敏感提示词——它是唯一一个
  // raw 永远不允许落取证文件的事件类型;今天没有 adapter 给它带 raw,
  // 这个守卫把脱敏保持为刻意的不变量而不是 adapter 行为的偶然
  if (event.raw !== undefined && event.type !== 'input.requested') {
    await this.writeRawBestEffort(runId, event)
  }
  await this.persistAndPublish(event, runId, this.clock.now())
}

private async writeRawBestEffort(runId: string, event: AgentProviderEvent): Promise<void> {
  try {
    await this.rawWriter.append(this.sessionId, runId, rawLineForEvent(event))
  } catch (err) {
    // 双保险:writer 内部已经吞错,这里再兜一层防「行为不端的 writer 实现」
    logger.warn({ err, sessionId: this.sessionId, runId, type: event.type },
      'best-effort agent event raw write threw at the emit seam')
  }
}

为什么 input.requested 特殊:旧编码器对这个事件类型显式写 rawJson: null(注释写明请求体可能含敏感提示词)。raw 改道文件后,编码器管不到文件路径了,脱敏必须搬到新的分叉点上——这是第 8 节第一条偏差的来源。

A.2per-run 串行链与 rawSeq 发号

并发场景真实存在:工具输出合并器的后台定时器会在 provider 事件 emit 的同时对同一个 run 发起 append。两个并发首写如果都去读文件行数发号,会拿到同一个 rawSeq。解法是每个 (sessionId, runId) 维护一条 promise 链,后来的 append 排在前一个后面:

apps/daemon/src/services/agent-event-file-writer.ts真实代码(节选)
const rawSeqByRun = new Map<string, number>()
const appendChainByRun = new Map<string, Promise<void>>()

return {
  async append(sessionId, runId, line) {
    const key = `${sessionId}/${runId}`
    const prior = appendChainByRun.get(key) ?? Promise.resolve()
    const next = prior.then(() =>
      appendRawLine(baseDir, rawSeqByRun, sessionId, runId, key, line),
    )
    // 链上某环 reject 也要保住链——后续 append 仍排在它后面
    appendChainByRun.set(key, next.then(() => undefined, () => undefined))
    await next
  },
  ...
}

发号本身分两种情况:进程内已有计数器就自增;进程重启后第一次写,读一次现有文件的行数来续号——同一次读还顺带检测出下一跳要讲的「开尾」:

apps/daemon/src/services/agent-event-file-writer.ts真实代码(节选)
// 一次读取喂两个事实:行数续 rawSeq;开尾(崩溃留下的无换行残行)告诉调用方先封口
async function readExistingTail(filePath: string): Promise<{ lines: number; openTail: boolean }> {
  try {
    const contents = await readFile(filePath, 'utf8')
    return {
      lines: contents.split('\n').filter((line) => line.length > 0).length,
      openTail: contents.length > 0 && !contents.endsWith('\n'),
    }
  } catch {
    return { lines: 0, openTail: false }
  }
}

A.3崩溃截断的尾行封口

场景:进程在 appendFile 写到一半时被杀,文件尾部留下半行 JSON、没有换行。下一次 append 如果直接写,新行会在残行后面——读取时两行拼成一段解析不了的字符串,一起被丢掉。修复是在重启后的首次 append 检测到开尾时,先补一个换行把残行「封口」,残行自己作为一条坏行被读取侧跳过,新行完好存活:

apps/daemon/src/services/agent-event-file-writer.ts真实代码(节选)
async function appendRawLine(baseDir, rawSeqByRun, sessionId, runId, key, line) {
  try {
    // 路径解析放在 try 里:坏 id 段和其他 append 故障一样吞掉——emit 在 DB 写之前
    // 调 append 且自身无 catch,这里抛出去就会丢掉已提交事件和广播
    const { sessionDir, filePath } = resolveRawLogPath(baseDir, sessionId, runId)
    await mkdir(sessionDir, { recursive: true })
    const { rawSeq, sealTail } = await nextRawSeq(rawSeqByRun, key, filePath)
    // 开尾必须先补上缺失的换行,否则本行粘上残段、读取时两行一起丢
    const separator = sealTail ? '\n' : ''
    await appendFile(filePath,
      `${separator}${JSON.stringify({ rawSeq, ...definedLineFields(line) })}\n`)
  } catch (err) {
    logger.warn({ err, sessionId, runId, type: line.type },
      'best-effort agent event raw append failed')
  }
}

路径安全也在这一层:sessionId / runId 作为路径段先过白名单式检查(拒绝 ..、分隔符、NUL),拼出的路径再做一次「必须在 baseDir 之内」的 resolve 校验。append 路径上校验失败被吞掉(同上,不能拖垮 DB 写);purge 路径上校验失败会抛——递归 rm 绝不允许拿着可疑路径跑。

A.4读取侧:为「可能已损坏」而生的读

取证文件被读的时机恰恰是它可能坏掉的时机(崩溃、手改)。所以读取对三类坏行全部容忍:JSON 解析失败、语法合法但形状不对(null、数组、rawSeq 不是数字),都跳过并计数告警,绝不让一次取证读崩掉:

apps/daemon/src/services/agent-event-file-writer.ts真实代码(节选)
export async function readEventRaw(baseDir, sessionId, runId): Promise<RawLine[]> {
  ...
  const lines = contents
    .split('\n')
    .filter((line) => line.length > 0)
    .flatMap((line) => {
      let parsed: unknown
      try { parsed = JSON.parse(line) } catch { skipped += 1; return [] }
      // 一条 null 行如果放进去,下面 sort 取 .rawSeq 就直接抛了
      if (!isRawLine(parsed)) { skipped += 1; return [] }
      return [parsed]
    })
    .sort((left, right) => left.rawSeq - right.rawSeq)
  if (skipped > 0) logger.warn({ skipped, sessionId, runId }, 'skipped malformed agent event raw lines')
  return lines
}

readEventRaw 目前是内部取证工具——没有任何 tRPC / CLI 入口调用它。这是刻意的:没有 UI 消费方之前不开对外面。

排查路标 · 旅程 A
症状从哪下手
取证文件缺行 / 有洞,但 DB 事件齐全agent-event-file-writer.ts:搜日志 raw append failed / raw write threw——某次 append 被 best-effort 吞了
rawSeq 重复或乱序agent-event-file-writer.tsappend 的 appendChainByRun 串行链——检查是否有调用绕过了 port 直写文件
读出来的行数比文件里少readEventRaw:日志 skipped malformed agent event raw lines;多半是崩溃残行或手改坏行
input.requested 的 raw 出现在文件里session-sink.tsemit 的类型守卫被动过——这是明确的不变量回归
事件流里读不到 raw 字段了符合预期:decodeEventRow 不再重建 raw,envelope 里本来也从没流出过它

6旅程 B:删除与孤儿清扫(#108)

文件不在 DB 里,FK 级联删不到它。这条旅程走「session 消失后文件怎么跟着消失」的三条路 + 一张兜底网。

全景 · 3 跳
硬删挂点 ×3
service.ts / use-cases
purge 排水后 rm
agent-event-file-writer.ts
启动孤儿清扫
index.ts

B.1三条硬删路径,同一个次序

session 可以被直接删,也可以随 task / project 级联删。三条路径的次序完全一致:先撕运行时(interrupt + dispose + 关闭事件订阅),再删 DB 行(权威成功边界),最后清盘上文件——上传目录和 raw 目录并排清,都是 best-effort:

apps/daemon/src/agent/service.ts真实代码(节选)
async deleteSession(sessionId: string): Promise<void> {
  await this.getSession(sessionId)
  await this.teardownSessionRuntimeBestEffort(sessionId, { terminate: false })
  await this.repo.hardDeleteSession(sessionId)
  // FK 级联只够得着 DB 行;不清盘上目录,blob 就成孤儿。行已删除,
  // 文件系统故障只记日志、绝不回滚删除
  await this.purgeSessionUploads(sessionId)
  await this.purgeSessionRawFiles(sessionId)   // 本 PR 新增,与上传清理并排
}

task / project 级联路径(use-cases 层)对每个被删 session 做同样三连;测试用 spy 钉住了三条路径都调用了 purgeSessionRawFiles

B.2purge 先排水再 rm

竞态:purge 跑 rm 时,per-run 串行链上可能还有排队的 append。如果 rm 先落地、排队的 append 后执行,它的 mkdir 会把刚删掉的目录复活,给已删除的 session 留下一坨僵尸字节,直到下次启动清扫才被收走。所以 purge 先等掉该 session 前缀的所有链尾,再 rm:

apps/daemon/src/services/agent-event-file-writer.ts真实代码(节选)
async purgeSession(sessionId) {
  const sessionDir = resolveSessionDir(baseDir, sessionId)   // 坏 id 这里直接抛,rm 不冒险
  // rm 前排空排队的 append:调用方先撕 sink 再 purge,不会有新 append 入队
  await Promise.allSettled(
    [...appendChainByRun.entries()]
      .filter(([key]) => key.startsWith(`${sessionId}/`))
      .map(([, tail]) => tail),
  )
  try {
    await rm(sessionDir, { recursive: true, force: true })
  } finally {
    // 无论 rm 成败,都把该 session 的计数器和链条目清掉,防内存泄漏
    for (const key of [...rawSeqByRun.keys()])
      if (key.startsWith(`${sessionId}/`)) rawSeqByRun.delete(key)
    for (const key of [...appendChainByRun.keys()])
      if (key.startsWith(`${sessionId}/`)) appendChainByRun.delete(key)
  }
}

B.3启动孤儿清扫:兜底网 + 归档例外

purge 是 best-effort,失败会留孤儿目录。兜底在 daemon 启动:扫 agent-events/ 下的一级目录,session 行还在 DB 就留,不在就整目录删。判活谓词就是「getSession 有没有返回行」——它不过滤归档(软删)session,所以归档 session 的取证文件在清扫中存活,这一点有专门测试钉着(防将来有人给 getSession 加归档过滤时静默误删归档数据):

apps/daemon/src/index.ts真实代码(节选)
const rawWriter = createAgentEventFileWriter({ baseDir: join(EYRIE_HOME, 'agent-events') })
const runnerManager = new RunnerManager(rawWriter)   // manager 把 writer 注入每个它拥有的 sink
...
await agentService.onDaemonStart()
try {
  await rawWriter.sweepOrphans(async (sessionId) =>
    Boolean(await agentRepo.getSession(sessionId)),   // 归档行也是 truthy → 目录保留
  )
} catch (err) {
  logger.warn({ err }, 'best-effort agent event raw orphan sweep failed during startup')
}
排查路标 · 旅程 B
症状从哪下手
删掉的 session 目录还在盘上日志 raw event purge failedservice.ts);重启 daemon 看 sweep 是否收走;还在就查 sweepOrphans 的判活谓词
归档 session 的取证文件消失了index.ts 的判活谓词 + agentRepo.getSession——若它开始过滤归档行,这就是事故根因
删除后目录被「复活」出一个新文件agent-event-file-writer.tspurgeSession 的排水逻辑——查是否有调用方在 purge 前没撕 sink,还在往链上排 append

7旅程 C:流式文本的合并落库(#109)

跟一串 message.delta 从 provider 到数据库:进缓冲、按什么条件出来、怎么保证和状态事件的先后次序、以及 run 终结时缓冲里的尾巴文本去哪了。这是两个 PR 里设计密度最高的部分。

全景 · 5 跳
emit delta 分支
session-sink.ts
分组缓冲
message-delta-coalescer.ts
边界/定时 flush
session-sink.ts
persistAndPublish(一行一事务) 终结点排水
service.ts / runner-manager.ts

C.1emit 的 delta 分支:raw 照旧逐条,DB 行改走缓冲

先看前后对比,再看代码。关键点:取证文件不受合并影响——仍然每个原始 provider chunk 一行;被合并的只有 DB 行和广播粒度。

以前(每 chunk)
emit(message.delta)
encodeEvent(含 rawJson)
appendEvent:一事务 + 两个计数器自增
broadcast 这一个 chunk
现在(每 chunk 进缓冲,flush 才落库)
emit(message.delta)
raw 有值 → 取证文件照旧写一行
coalescer.accept:文本进 (run,session,role,item) 缓冲
flush 条件满足(16KB / 250ms / 边界事件 / 换组)才 appendEvent 一行拼接文本
broadcast 合并后的块(直播从逐 token 变 ≤250ms 块状)
apps/daemon/src/agent/session-sink.ts真实代码(节选)
if (event.type === 'message.delta') {
  if (event.raw !== undefined) await this.writeRawBestEffort(runId, event)  // 取证逐条不动
  const deltas = this.messageDeltas.accept(runId, this.sessionId, event)
  for (const delta of deltas) {   // accept 可能同步吐出已满/换组的块,立刻落库
    await this.persistAndPublish(delta.event, delta.runId, delta.createdAt)
  }
  return   // delta 本身不直落——它的文本已在缓冲里
}

时间戳语义跟着变:合并行的 createdAt 是 flush 时刻(注入的 clock 给出,测试可控),occurredAt 是缓冲里第一个 chunk 的 provider 时间——一行代表一段时间窗的文本,两端各取其一。

C.2分组缓冲与单开组指针

accept 做三件事:必要时先关掉上一个组、把文本进组、检查字节阈值。第一件事最微妙——每条 (run, session) 时间线上只允许一个开组

apps/daemon/src/agent/message-delta-coalescer.ts真实代码(节选)
accept(runId, sessionId, event): RunScopedMessageDelta[] {
  const key = bufferKey(runId, sessionId, event.role, event.itemId)
  const flushed = this.flushGroupBoundary(runId, sessionId, key)  // 换组先关旧组
  const buffer = this.getBuffer(key, runId, sessionId, event)

  buffer.chunks.push(event.text)
  buffer.bytes += Buffer.byteLength(event.text, 'utf8')
  if (buffer.occurredAt === undefined && event.occurredAt !== undefined) {
    buffer.occurredAt = event.occurredAt   // 只记首个 provider 时间
  }

  this.openGroups.set(timelineKey(runId, sessionId), key)
  if (buffer.bytes >= this.byteThreshold) return [...flushed, ...this.flushGroup(key)]

  this.ensureTimer(key, buffer)   // 没到阈值就挂 250ms 定时器兜底
  return flushed
}

// 下一个 delta 换了 role 或 itemId 时关闭开组:回到早先组的 delta 会另起一行,
// 而不是跨过中间到达的组去合并——否则回放顺序被重排
private flushGroupBoundary(runId, sessionId, incomingKey) {
  const openKey = this.openGroups.get(timelineKey(runId, sessionId))
  if (!openKey || openKey === incomingKey) return []
  return this.flushGroup(openKey)
}

为什么需要单开组指针而不是各组独立缓冲到各自超时:设想到达序是 assistant(item-a) → reasoning(item-r) → assistant(item-a)。如果两段 item-a 合并成一行,中间到达的 reasoning 在回放里就被挤到它们之后——文本内容没错,但次序错了。切组即 flush 保证「flush 顺序 = 到达顺序」,这有专门测试钉住。

C.3非 delta 边界:全量排水,且排的是所有 run

任何非 delta 事件(审批请求、消息完成、run 完成、工具事件……)落库前,先把缓冲的文本全部排空——保证合并行拿到比边界事件更小的 sessionSeq,回放时文本仍在边界之前。这段源码注释值得原样细读,它讲清了两个非直觉决策:

apps/daemon/src/agent/session-sink.ts真实代码(节选)
// 在这个非 delta 边界前 flush 缓冲的 delta,让合并行拿到比后续事件更小的 sessionSeq。
// 排空的是【所有 run】的缓冲,不只这个事件自己的 run:sessionSeq 是 session 级全局的,
// 一个钉在旧 run 上的 resolution 事件会反超当前 run 上已接收的文本(那些文本会被
// 定时器以更大的 sessionSeq 提交,到达顺序就反转了)。这里的次序是 load-bearing 的,
// 依赖 appendEvent 在同步的 better-sqlite3 事务里分配 sessionSeq(读计数器和提交之间
// 没有 await 让出点):flush 完成并提交后,调用方才恢复去落边界事件。若换成一个
// 分配 sessionSeq 前会 await 的驱动,边界反超待 flush 文本的窗口就重新打开。
await this.flushMessageDeltas()

// 完成中的工具先排空它的合并输出(尾部 delta 保住更小的 sessionSeq)再落终结事件
if (event.type === 'tool.completed') await this.flushAndRetireTool(runId, event.toolCallId)

「排所有 run」是第二轮机器人评审补上的(首版只排边界事件自己钉的 run),第 8 节有完整偏差记录。

C.4定时器 flush:fire-and-forget,但丢失必须可观测

低流量文本到不了 16KB,靠 250ms 定时器兜底。定时器回调在任何 await 调用方之外触发,写库失败没有地方向上抛;而缓冲已经被排空,这块文本就永久丢了。设计的选择是:不重试,但用 error 级日志 + 字节数把丢失暴露出来:

apps/daemon/src/agent/session-sink.ts真实代码(节选)
this.messageDeltas.setFlushHandler(({ runId, event, createdAt }) => {
  // 定时器在任何被 await 的调用方之外触发,persist 失败无处上浮;缓冲已排空,
  // 写失败即这块合并文本永久丢失——记日志,不静默吞掉
  void this.persistAndPublish(event, runId, createdAt).catch((err) => {
    logger.error(
      { err, sessionId: this.sessionId, runId, role: event.role, itemId: event.itemId,
        textBytes: Buffer.byteLength(event.text, 'utf8') },
      'dropped a coalesced message.delta after a background flush write failed',
    )
  })
})

C.5绕过 sink 的终结点:谁来排最后一口水

缓冲引入了一类新风险:run 的终结是直接的 repo 写,不经过 sink。如果终结时缓冲里还有文本,定时器会在 run 已经 terminal 之后才把它落库——可能插进下一轮 turn 的行序里。所以每个终结点都要先排水。四个位点,语义各有讲究:

① 控制命令(runControl)——成功路径的 flush 是严格的:这次 flush 是命令期间流出的低于阈值文本唯一的落库机会,写失败必须让 run 记为 failed 并上抛,不能让 run 标着 completed 而输出被静默丢掉。失败路径的 flush 则是 best-effort——run 本来就在失败,flush 的错不能遮住原始错误:

apps/daemon/src/agent/service.ts真实代码(节选)
try {
  const result = await handle.runner.runControl(command)
  await handle.sink.flushMessageDeltas({ runId: run.id })   // 直通,不包 try——失败走 catch
  await this.repo.finishControlRun(run.id, result.status === 'error' ? 'failed' : 'completed')
  return result
} catch (err) {
  await this.flushLiveMessageDeltas(sessionId, handle.sink, { runId: run.id })  // best-effort
  await this.repo.finishControlRun(run.id, 'failed')
  throw err
}

② 用户打断(interruptCurrentRun)——排两次。第一次在 interrupt() 之前:interrupt 的 await 可能让 runner 的 closed promise 先落定,crash 分支会销毁 sink(缓冲直接作废),等 interrupt 返回时文本已经没了。第二次在 terminateRun 之前:interrupt 悬停期间 provider 还可能流进新文本:

apps/daemon/src/agent/service.ts真实代码(节选)
const handle = this.runnerManager.get(sessionId)
if (handle) {
  // 在 await interrupt 之前排水:interrupt 悬停时 closed promise 可能落定、
  // crash 分支销毁 sink 丢掉缓冲——先 flush 保住用户主动停止前的尾部文本
  await this.flushLiveMessageDeltas(sessionId, handle.sink, { runId: run.id })
  try { await handle.runner.interrupt() } catch { /* best-effort,下面的 CAS 才是权威 */ }
  // 再排一次:interrupt 期间新进的文本;若 crash 路径已销毁 sink,这里是无害 no-op
  await this.flushLiveMessageDeltas(sessionId, handle.sink, { runId: run.id })
}
await this.repo.terminateRun(run.id, 'interrupted')

③ 优雅销毁(runner-manager dispose,close/delete/archive 都走这里)——dispose 前排水;④ 崩溃——刻意不排。崩溃路径直接 destroy() 丢弃缓冲,最多损失一个 flush 窗口(≤250ms)的文本,这是锁定的取舍(有测试钉住「崩溃不落缓冲行」):

apps/daemon/src/agent/runner-manager.ts真实代码(节选)
async dispose(sessionId: string): Promise<void> {
  ...
  this.disposing.add(sessionId)
  await this.flushSinkMessageDeltas(sessionId, entry.sink)  // 优雅路径:先排水(best-effort)
  await entry.runner.dispose()
  ...
}

private deleteEntry(sessionId: string): void {
  this.runners.delete(sessionId)
  // 先 destroy 再放引用:定时器不能再往一个 manager 已不拥有的 session 上开火
  this.sinks.get(sessionId)?.destroy()   // 崩溃路径经 handleRunnerExit 到这——不 flush,直接弃
  this.sinks.delete(sessionId)
}

还有一个容易漏看的小位点:startTurn 在 runner 尚不存在时会临时造一个 fallback sink 来发 message.user 事件。现在 sink 内部持有真实的 interval 定时器,这个「用完即弃」的 sink 必须在 finally 里 destroy(),否则定时器泄漏。它只发非 delta 事件、永远不会有缓冲,destroy 纯粹是拆定时器。

排查路标 · 旅程 C
症状从哪下手
回放里一段文本出现在审批/工具事件之后(到达时明明在前)session-sink.tsemit 非 delta 分支的 flushMessageDeltas()(无参=全 run 排水)是否被改窄;appendEvent 是否仍是同步事务分配 sessionSeq
UI 直播文本一段一段跳而不是逐字符合预期:直播粒度=flush 窗口(默认 250ms/16KB),调 messageDeltaCoalescer options 可改
一段回复的结尾几个字丢了(run 被打断/关闭后)service.ts:对应终结点的排水调用;日志 message delta flush failed;若是进程崩溃,≤250ms 窗口内的损失是设计内行为
日志出现 dropped a coalesced message.delta定时器后台 flush 写库失败,该块文本已丢(textBytes 字段给出规模)——查 DB 写失败的根因
控制命令 run 标 completed 但输出缺失service.ts runControl:成功路径 flush 必须直通(失败要 finishControlRun('failed'))——若被包成 best-effort 即回归
两段不同 role/itemId 的文本被拼进一行message-delta-coalescer.tsbufferKeyflushGroupBoundary 的 openGroups 指针

8设计意图 vs 落地偏差

实现前有一份锁定了各项决策的设计稿,实现后又经历了多轮交叉评审 + PR 机器人评审。照稿做成的部分上面都讲过了;下面是实际落地和最初设计不一样的地方——每条都是一道认知裂缝,也是这两个 PR 里最值得记住的教训。

设计稿说实际做成为什么变
raw 写入守卫字面只有 event.raw !== undefined 额外排除 input.requested;事后被确认为正式锁定的不变量 旧编码器对该事件显式置空 raw(请求体可能含敏感提示词)。raw 改走文件后绕过了编码器的置空,不补守卫脱敏就静默失效
rawSeq = per-run 计数器,首写读文件行数续号 外面又包了一层 per-(session,run) promise 串行链 评审发现两个并发首写会同时读到相同行数、发出重复 rawSeq(工具输出定时器与 provider 事件并发是真实场景)
(未提崩溃残行) 首写检测「开尾」,先补换行封口再 append 机器人评审指出崩溃截断的残行没有换行,下一条 append 会粘上去,读取时两行一起丢
purgeSession = rm 目录 rm 前先 Promise.allSettled 排空该 session 的所有 append 链尾 机器人评审指出已入队的 append 在 rm 后执行会经 mkdir 复活目录,给已删 session 留僵尸字节
合并按 (run,session,role,itemId) 分组即可 加了 openGroups 单开组指针:换组即 flush 旧组 评审发现回到早先分组的 delta 会跨越中间到达的组合并,回放次序被重排(内容对、顺序错)
非 delta 边界 flush 事件所在 run 的缓冲 改为 flush 所有 run 的缓冲 机器人二轮发现:钉在旧 run 的 resolution 事件会先拿 sessionSeq,当前 run 已缓冲的文本之后才由定时器落库——到达序反转。sessionSeq 是 session 级全局的,边界排水的作用域也必须是
(未列 runControl / interrupt 这类不经 sink 的终结点) 三处终结点补 drain;runControl 成功路径的 flush 失败会把 run 记 failed 并上抛(其余 drain 是 best-effort) 机器人评审发现 run 被直接 repo 写终结后,缓冲文本会在 terminal 之后由定时器落库、可能插进下一轮的行序;其中成功路径最初被做成 best-effort,二轮又指出「flush 吞错会让 run 标 completed 而输出丢失」
两处评审提了但刻意不改:崩溃路径不强制 flush(损失 ≤ 一个 flush 窗口是锁定取舍,强 flush 会让 crash 清理依赖一次可能失败的 DB 写);迁移基线时间戳重置(开发期单基线策略:schema 变更即重建本地库,无已发布用户)。
一条元教训:后三条偏差是同一族问题的三次发现——「缓冲引入后,哪些写路径绕过了 sink、边界该覆盖哪些缓冲」。第一次修了 runControl,第二次修了 interrupt,第三次才把边界 flush 的作用域修全。将来给任何写路径加缓冲时,值得一次性枚举所有终结点和边界作用域,而不是等评审逐个找。

9心智模型补丁

agent_events 每行带 raw_json,想看原始 provider 载荷查表即可 表里没有 raw 了;原始载荷在 $EYRIE_HOME/agent-events/<sessionId>/<runId>.raw.jsonl,用 readEventRaw(纯内部 helper)读
一个 provider message.delta chunk = agent_events 一行 一行 = 一个 flush 窗口内同 (run,role,itemId) 的拼接文本;取证文件仍是每原始 chunk 一行
数 DB 行数≠数 provider chunk 数;对账要拿取证文件。
订阅者按 provider 发出的粒度逐 chunk 收到直播文本 直播粒度 = flush 窗口(默认 250ms / 16KB),块状到达;平滑显示是客户端的事
事件行 createdAt ≈ provider 发出时刻(emit 时的 Date.now()) 合并行 createdAt = flush 时刻(注入 clock),occurredAt = 窗口内首个 chunk 的 provider 时间;非合并事件不变
emit 一旦 resolve,这条事件的所有痕迹都已持久 delta 的 emit resolve 只代表进了缓冲;持久要等 flush。优雅终结(close/interrupt/delete/控制命令)会先排水,进程崩溃则丢弃缓冲(≤1 个窗口)
删 session = DB 级联搞定一切 盘上还有两处要清:上传目录 + raw 事件目录,都在 DB 删除提交后 best-effort 清;漏网的靠 daemon 启动孤儿清扫兜底;归档 session 的文件刻意保留
SessionSink 是个轻对象,new 出来用完丢掉即可 sink 内部持有真实 interval 定时器(两个 coalescer),临时 sink 必须 destroy(),否则定时器泄漏——startTurn 的 fallback sink 就是这么处理的

10新词表

取证文件域(#108)
forensic raw fileper-run 的 append-only JSONL,存 adapter 发出的原始 provider 载荷;只为事后排查,从不进事件流
rawSeq取证文件内的行序号,per-run 自增,与 DB 事件 seq 无关
append chainper-(session,run) 的 promise 链,串行化并发 append,保 rawSeq 单调且盘上顺序一致
open tail / seal崩溃留下的无换行残行叫开尾;下次首写先补换行「封口」,让残行独立成一条可跳过的坏行
orphan sweepdaemon 启动时扫 agent-events/ 一级目录,session 行已不存在的整目录删除
合并域(#109)
MessageDeltaCoalescermessage.delta 的内存缓冲合并器,镜像既有的工具输出合并器(注入 clock + timer,可测)
group / bufferKey合并的最小单位,(runId, sessionId, role, itemId) 四元组;不同组的文本永不拼接
open group每条 (run,session) 时间线上唯一「还在收文本」的组;来了不同组的 delta 就先关它——保 flush 序=到达序
boundary flush任何非 delta 事件落库前,先排空所有 run 的缓冲,保证合并行的 sessionSeq 小于边界事件
flush window时间兜底阈值(默认 250ms);也是崩溃时文本损失的上界
fallback sinkstartTurn 在 runner 尚不存在时临时创建、发完 message.user 即 destroy 的一次性 sink

11测试与风险地图

有兜底的(测试钉住的行为)

薄冰(无测试或已知遗留,纯事实陈述)

合并前必办:无。两个 PR 的 CI 与本地全量门禁(typecheck / lint / check / test / build)均绿,机器人评审全部闭环(每条要么修掉要么带 by-design 回复)。次序约束一条:先合 #108,#109 是它的 stacked 分支。

12验收提示

以下几处容易在验收时被误判成缺陷,实际都是刻意为之:

13覆盖声明

本报告作者全量逐行精读了两个 PR 的完整 diff(#108:+854/−104;#109:+1258/−36,含全部测试代码),未抽样、未略读、未使用二手转述。报告中所有代码片段均出自作者亲自读过的 head 分支源文件(session-sink.tsmessage-delta-coalescer.tsagent-event-file-writer.tsservice.tsrunner-manager.tsindex.ts)与 base 侧 session-sink.ts 原文,亲手裁剪。偏差章节的素材来自实现过程记录与 PR 评审线程,结论已全部内联。