系统设计 / 任务运行时 · 从一次 run() 到可观测的调度层 / 观测与调试:事件流与快照 待审核 8 / 10
事件流 · devtools

观测与调试:事件流与快照

调度层的状态几乎都是瞬时的:一个任务从 queued 到 running 只隔几毫秒,取消传播完就不再留痕。等到有人报告「列表偶尔加载不出来」,那一批任务早已离场。@vega/job 为此留了三个观测口:一条按发生顺序投递的事件流、一份只回答「此刻」的快照、一块把两者画出来的 devtools 面板。本页把这三个口各看一遍,再落到运行时自己的收尾动作上。

1 · 事件流的形状

subscribe(runtime, listener) 注册一个 listener 并返回退订函数,运行时每推进一步就向所有 listener 投一条事件。事件是三段式的:type 是事件名,at 是时间戳,payload 装这条事件特有的字段。任务 id 在 payload.id 里,不在顶层。

一个固定场景可以把这条流数清楚。父任务派出两个子任务、三个都正常完成,共 12 条事件;每个任务各自经过 task:registeredtask:queuedtask:startedtask:done,四条一组,父任务的登记最先到、完成最后到(lab/facts.test.ts §1)。事件条数与任务数的这种整齐关系只在顺利路径上成立:取消会插入 task:cancellingafter 依赖会插入 task:parked,重试与进度另有自己的类型。

listener 拿到的静态类型是 JobEvent{ readonly type: string; readonly at: number; readonly payload?: unknown }——payloadunknown,编译器不允许直接读 ev.payload.id。包里另有一个 BuiltinJobEvent,把内建的每种 type 与它对应的 payload 配成可辨识联合,switch (ev.type) 之后每个分支的 payload 就窄成了具体形状。

由此引出一个绕不开的取舍:总线是开放的,扩展层(coordination:*debounce:*)和调用方自己发的事件走同一条通道,所以消费方的类型一旦写成 BuiltinJobEvent | JobEvent 去兜住未知事件,联合里那个 JobEvent 就把 narrowing 整个抹平了——本系列的 core/model.ts 正是这么声明的,代价就是 ingest 里读 payload 全靠手写断言。要么放弃兜底换 narrowing,要么兜底并自己承担断言,两者只能选一个。

无论选哪个,顶层永远只有 type / at 两个字段:任务 id 一律在 payload.id。照着「ev.id」写出来的 listener 会在 undefined 上取值抛错,而且抛出之后未必看得见,见下面这条。

警示 · listener 抛错不会打断调度:错误被交给 hooks.onIsolatedError,scope 为 'listener',meta 里带 source: 'subscribe()' 与出事那条事件的 eventType。同一台运行时上的其他 listener 照常收到全部事件,任务也照常跑到 successlab/facts.test.ts §1)。反过来说,构造运行时时没传 hooks.onIsolatedError 就等于把这类错误丢进真空:listener 写挂了,页面上一行报错都不会有,事件日志只是静悄悄地空着。onIsolatedError 的签名是单参的 report 对象 { scope, error, meta }meta 的形状随 scope 变(见 engine/hooks.tsIsolatedErrorReport)。

把事件按用途分档是消费方常见的需求:task:registered 这类稳定的生命周期事件是一档,task:progressdebounce:*coordination:* 这类诊断事件是另一档。运行时不提供这层过滤——subscribe(listener) 就一个参数,没有 category 选项,包里也没有导出任何判定函数。分档只能由订阅者自己按 type 前缀做,本页 lab 的 observe.store.ts 里那张 DEBUG_TYPES 集合就是一份可以照抄的映射。

subscribe() 的另一个限制是只收得到注册之后的事件。晚挂载的面板要补历史,得在建运行时时开 events.retain:总线按这个尺寸维护一个 ring buffer,getRecentEvents() 一次性返回其中的内容。不开 retain 时它恒为空数组;开到 5 时返回的恰是流的末尾 5 条,开到 100 时 12 条全在(lab/facts.test.ts §1)。retain 是硬上限,写满即丢最旧的一条,没有按时间的过期。

图 1-1 · 七个场景在真实运行时上跑出的事件流,左侧按 type 计数,日志每行标出它属于哪一档。可切换档位过滤,或调 events.retain 后重跑父子场景看 getRecentEvents 回填几条。

注 · 事件里的 at 与 payload 里的 queuedAt / startedAt / endedAt 取自同一次 Date.now() 调用(engine/task-observer.ts),画时间轴用哪个都一样。包里另外导出了一个 now()performance.now() 的封装),它不受系统时钟调整影响,更适合算时间差;但事件流本身用的是 Date.now(),两者混着减会得到无意义的值。本系列所有 Gantt 都只用事件里的这些 Date.now() 时间戳。

2 · inspect() 的此刻

事件流回答「发生过什么」,inspect() 回答「现在有什么」。它返回五份数据:queue 是已排队未开跑的任务,running 是在跑的,parked 是还没进队列的(等 after 依赖,或被同 key coordinator 扣住),tree 是 id 到父子关系的映射,deps 是 id 到它 after 的那些 id。

三份任务清单里放的是 ReadonlyTask 而不是 Task:有 statusstategetSnapshot()subscribe(),没有 cancel()retry()。这是故意的收窄,inspect() 不该被当成控制面板用;要控制某个任务,得握着 run() 当初返回的句柄,或者用 getTaskInstance(id) 按 id 取回活的实例。

快照只覆盖活着的任务。一个任务抵达终态就退出这三份清单,tree 随之清空,getTaskInstance(id) 从此返回 undefinedlab/facts.test.ts §2)。所以两个观测口是互补的:快照给出结构,事件流给出历史,devtools 面板两样都要订。

图 2-1 · 一台常驻运行时的 inspect() 快照,三份清单随事件刷新,下方是 tree 与 deps。可派任务、暂停调度或 cancelAll,单击任一行调 getTaskInstance(id) 看它是否还在图里。

depstree 是两张不同的图。tree 记的是谁派出了谁,决定取消如何级联(见任务与状态机 §4);deps 记的是谁在等谁,决定排队顺序(见编排与依赖)。同一个任务可以在两张图上各有一组邻居。

3 · devtools 面板

mountTaskDevtools(container, runtime)@vega/job/devtools 导入,往容器里挂一块实时面板,返回卸载函数。面板自己 subscribe() 一份事件流、用 requestAnimationFrame 合并重绘,分区呈现:在飞任务与队列并排在顶部,其下依次是错误、dep-stuck、coordination 决策、被拒的调用、诊断计数与等待统计,最后是任务树与时间线。它自带一套深色样式,<style> 元素注入 document.head,卸载函数会连同面板 DOM 一起移除。

import { mountTaskDevtools } from '@vega/job/devtools';

const stop = mountTaskDevtools(host.value!, runtime);
onUnmounted(() => { stop(); void disposeRuntime(runtime); });

面板基本是只读的,唯一的例外是错误区:失败且 retryable 的任务上挂着一个 retry 控件,按下去调的是 getTaskInstance(id)?.rerun()。事件 payload 必须可序列化、塞不进闭包,所以 task:done 只带一个 retryable 布尔,重试入口得由面板自己按 id 取实例。取实例的时机很紧:事件派发是同步的,task:done 的 listener 一返回,那个节点就从图里注销了(见 §2),所以面板必须在 handler 里当场把实例抓在闭包中,getTaskInstance 也正为此做成同步 API。

图 3-1 · 库自带的 devtools 面板挂在一台并发 4 的运行时上,八个按钮各触发一种典型形态。可连点 parallel burst 与 priority queue 观察队列区重排,或触发 error task 后在面板里按 retry。

这八个场景移植自 example/job/devtools.html,但那份 demo 停在旧的调用形态:优先级写在 policy 里、重试写成 RunOptions.retry、去抖写成 policy.debounceMs。这些字段现在都不存在,运行时会原样忽略它们,页面上仍然「有反应」,只是反应的是默认档位。本页把它们改写成了 priority: TaskPriority.*TaskErrorretryable 标记与 debounce() shaper。

另一个面板 mountReactiveDevtools(container, graph) 挂的不是运行时而是响应式图。图上没有事件总线可订阅,面板走 graph.onChange 重渲染,详见响应式依赖图

4 · 标签与关联

任务在事件流里靠三种标识被认出来。labelRunOptions 里的一个字符串,不参与任何调度,只为让面板与日志上出现「upload」而不是 t17flowIdstepIdx 在任务由 flow() 派出时自动带上,把一串任务归到同一次编排下。correlationId 则跨越重试链:task.rerun() 派出的是一个全新实例,id 换了,correlationId 仍是最初那个(lab/facts.test.ts §4),追踪层据此把几次尝试缝回同一次逻辑调用。

hooks 是与事件流并行的另一条通道,在建运行时时传入:onTaskStart / onTaskSuccess / onTaskError / onTaskCancelled / onTaskProgress 各对应一次状态迁移,onIsolatedError 收所有被吞掉的错误。它与事件流的分工不在信息量,在粒度。progress 尤其明显:taskFn 在同一个同步段里连调 5 次 ctx.progress(),hook 触发 5 次,而总线上只有 1 条 task:progresslab/facts.test.ts §4)。事件总线按 task id 做 microtask 合并,最新值胜出,为的是不让紧循环里的进度洪水淹掉 devtools。需要每一次增量就订 hook;hook 里又要更新反应式状态时,用 keyedThrottle(fn, ms) 按 id 节流,它保留最后一次值,不会把终值丢掉。

effects 的清理失败也走 onIsolatedError,scope 是 'cleanup'tracingEffecttraceContext 会被展开进这条报告的 meta,与 effectNamesource: 'effect.dispose' 并列,于是一次 span 收尾时的异常能带着自家的 trace id 落到可观测管道里(effects 本身见作用域与 effects)。

5 · 运行时的收尾

whenIdle()dispose() 都返回 Promise,都在「一切都结束了」时 resolve,但它们对在飞任务做的事正相反。

whenIdle() 是纯观察点:它快照调用那一刻的待处理工作,等这批任务连同它们的 effect 屏障全部落终态。调用之后新派的任务不会延长等待,运行时也不会被封口。三个在飞任务在 whenIdle() 之后仍然全部 success,随后 run() 照常受理(lab/facts.test.ts §5)。

dispose() 则先取消再等待:它取消活跃、排队与 parked 的任务,清掉 debounce 与 throttle 的定时器,然后等每个任务的 taskFn 体与异步 effect 收尾跑完。被它取消的在飞任务拿到的 cancelReasonqueue-clear 而不是 abort,因为这条路径走的是 cancelAll() 的结构性取消。dispose 之后 run() 会立刻失败,那次调用在总线上只留一条 task:rejected(reason 为 runtime-disposed),既没有 task:registered 也没有 task:donetask:rejected 是唯一一种不隶属于某个已登记任务的生命周期事件,按 id 建索引的消费方要单独处理它——payload.id 指向的那个 task 从来没进过图。

排空结束时发一条 runtime:disposed,这是整条事件流唯一的收尾 sentinel,订阅者可以把它当作「现在可以安全重置状态机」的信号。它恰好一条,重复调用 dispose() 共享同一次排空(lab/facts.test.ts §5)。排空卡住时另有两条诊断事件:warnAfter 到点发 runtime:dispose-stalled 并继续等,timeout 到点发 runtime:dispose-forced 后立刻 resolve,此时 runtime:disposedforced: true,未结束的屏障被有意泄漏。

图 5-1 · 三个在飞任务分别被 whenIdle() 与 dispose() 收尾,时间轴末尾那条是收尾之后补发的第四个任务。可切换两种收尾方式,对照三个任务的终态与 runtime:disposed 的条数。

建议 · runtime 是不透明 handle,没有方法,因此 Symbol.asyncDispose 挂不上去,await using 对它不适用(createScope 返回的作用域仍是带方法的对象,那里照旧可以用)。脚本或测试里显式收尾:try { … } finally { await disposeRuntime(runtime) }。组件里则相反:onUnmounted 里不必 awaitvoid disposeRuntime(runtime) 让它自己排空即可,否则卸载会被在飞任务拖住。

观测口到此为止:事件流、快照、面板、hooks,加上两个收尾动作。下一页把视角从「任务」抬到「值」,看建在同一台运行时上的响应式依赖图