职责与边界
Worker 负责:- 对外提供 gRPC
WorkerService(api/grpc/worker/worker.proto):RequestTaskSlot、SubmitTask、GetTaskResult、WatchTasks。WatchTasks是 server-streaming RPC,Worker 只监听一个 TCP 端口。协议细节见 协议与接口。 - 维护本地 slot 表和 task 索引,做原子准入和
taskId级幂等。 - 为每个 task 固定一个
taskCodeEpoch,串行跑 Call Builder,把 call list 交给 dispatcher 并行执行,最后串行跑 Writer Plugin。 - 收敛
TaskResult,推送给WatchTasks订阅者,并按resultRetentionMs保留。 - 向 etcd 注册并周期上报心跳,供 Coordinator 选址。
- 生成
taskId、构造 payload、DAG 编排或定时触发(这些在上游)。 - 跨 Worker 的去重、全局结果查询或崩溃后的本地恢复。重启后旧 slot、旧结果和订阅全部丢失。
- call 级调度细节(dispatcher、executor adapter)。这部分见 Call 执行子系统。
WorkerCore是纯内存、Sans-IO 的决策引擎。它不读时钟,每个方法都显式接收now。- slot 表是容量和 task 状态迁移的唯一真相;etcd 心跳只是它的派生视图。
slotTable与taskIndex在同一次方法调用里一起更新,不存在撕裂状态。- 一旦 task 成功激活,后续所有失败(Builder、call、Plugin、超时、Watch 断连)都收敛为
TaskResult里的终态,不再返回提交层错误。 - adapter 不维护第二份 slot 或 phase 真相。它只执行 core 输出的命令,再把结果作为事件回送。
代码位置
核心类型与接口
core(internal/worker/core/)
WorkerCore(worker.go):单线程状态机。持有slots []SlotEntry、taskIndex map[string]*TaskIndexEntry、tasks map[string]*TaskContext、ready、draining。- 输入事件(都在
worker.go,除非注明):HandleRequestTaskSlot(taskID, ttlMs, now)、HandleSubmitTask(req, adapterTaskCtx, now)、HandleTaskActivationPrepared(taskID, slotID, taskCodeEpoch, now)、HandleTaskActivationFailed(taskID, failureCode, now)、HandleTaskPhaseFinished(taskID, phase, outcome, now)、HandleTaskTimedOut(taskID, now)、HandleWatchAttached/HandleWatchDetached/HandleWatchDisconnected(taskID, now)、HandleSlotLeaseExpired(slotID, now)(worker_slots.go)、HandleGetTaskResult(taskID)。 - 定时扫描辅助:
ExpiredSlots(now)、TimedOutTasks(now)、WatchDisconnectedTasks(now)、PurgeExpiredResults(now)。adapter 每个 tick 调用它们,再把结果喂回对应Handle*。 - 输出命令(
types.go):PrepareTaskActivation、StartBuilderPhase、StartCallPhase{CallList, Streaming}、StartWriterPhase{Outputs}、ConvergeTask{State, Result}。它们嵌在结果类型RequestTaskSlotResult、SubmitTaskResult、ActivationResult、PhaseTransitionResult、LeaseExpiredResult、GetTaskResultOutput里返回。 TaskContext(types.go):TaskID、SlotID、TaskCodeEpoch、Phase、DeadlineMs、WatchAttached/WatchEverAttached/WatchDisconnectDeadlineMs、CallFailurePolicy、CallRetryPolicy、AcceptedAtMs、AdapterTaskCtx,以及三个阶段的PhaseOutcome。DeriveHeartbeat(now) etcd.WorkerHeartbeat:从 slot 表派生ReservedSlots/RunningTasks。
internal/worker/adapters/)
GRPCServer(grpc_server.go):实现RequestTaskSlot/SubmitTask/GetTaskResult,把WatchTasks委托给SubscriptionManager。workerErrToStatus(wire.go)把core.WorkerError映射为 gRPC status,业务错误码经workerCodeToGRPC映射:InvalidArgument → InvalidArgument、NoSlot → ResourceExhausted、Unavailable → Unavailable、InvalidSlot → FailedPrecondition、NotFound → NotFound、SlotUncertain → Aborted。Orchestrator(orchestrator.go):桥接WorkerCore和DispatcherCore。入口方法RequestTaskSlot、SubmitTask、GetTaskResult、TickTimers、DeriveHeartbeat、SetDraining、WaitDrain、RunDispatchLoop、SyncExecutorSnapshots、CheckWedgedExecutors。OrchestratorDeps(orchestrator.go):构造依赖,包括BuilderRegistry、PluginRegistry、IOScope、Admission、FunctionCodeViewProvider、EpochResolver、ExecAdapter、AuditGate、StreamBuild。taskRuntime(task_runtime.go):adapter 侧的 task 运行时记录,持有ioTaskCtx、TaskIOScope、ctx/cancel、functionCodeView、流式生产句柄和TaskStatistic。它不维护 phase 或 slot 真相。SubscriptionManager(stream_subscriber.go):每条WatchTasks流一个streamSubscriber;把首条 attach / 最后一条 detach 翻译为Orchestrator.HandleWatchAttached/HandleWatchDetached。EtcdPublisher(etcd_publisher.go):Register、RunHeartbeatLoop、Deregister。key 为KeyPrefix + workerAddr,lease TTL 默认 10s。ExecutorAdapter接口(exec_adapter.go):Orchestrator 依赖的 UDS executor 面,生产实现是executor.Adapter。
internal/worker/app/)
Profile(app.go):Deployment(block/bundle,只做观测标识)、Service(tracing / Prometheus 服务名)、WorkerRegistryPrefix(etcd 注册前缀)、UsageService、Builders、Plugins、TuneDefaults(只改内置默认值)。Run(p Profile):唯一的装配和启停序。
数据流 / 执行流程
Worker 采用”输入事件 → core 决策 → adapter 执行命令 → 事件回送”的模式。Orchestrator 用 o.mu 保护 WorkerCore,用 runtimeMu 保护 runtimes;DispatcherCore 只在 dispatch actor(RunDispatchLoop goroutine)上被修改,其他 goroutine 通过 postDispatchEvent 投递事件。锁序和 actor 不变量写在 orchestrator_actor.go 文件头注释,并由 actor_invariants_test.go 强制。
几个关键点:
SubmitTask在返回前只做到HandleSubmitTask。core 立即创建TaskContext(Phase = Preparing),使重复提交命中幂等路径;激活和阶段执行在activateAndRungoroutine 中异步进行。- 激活失败(epoch 解析失败、
Pin失败、NewTaskScope失败)走HandleTaskActivationFailed,失败码为ACTIVATION_FAILED或IO_SCOPE_FAILED。 - Builder 阶段:
runBuilderPhase先取builderSem,再从BuilderRegistry查FunctionCallConfig.Type。查不到直接以BUILDER_NOT_FOUND收敛。若 builder 实现了StreamingCallBuilder且StreamBuild.Enabled,Builder 阶段只跑PrepareStream,回送CallStreamMarker,core 输出StartCallPhase{Streaming: true},随后runCallStreamPhase在scanSem下启动生产 goroutine。 - Calls 阶段:
runCallPhase取executorSem,在锁外用dispCore.PrepareCallPhaseWithDigests准备 call 上下文,再投递evCallPhaseStarted给 actor。call 级调度、重试、subcall 归 Call 执行子系统。actor 收到ConvergeTaskCalls命令时调用HandleTaskPhaseFinished(PhaseCalls)。 - Plugin 阶段:
runWriterPhase取pluginSem,按ResultHandler.Type查PluginRegistry。未配置ResultHandler时不执行任何插件,直接以成功回送。 - 终态:
executePhaseResult对仍 in-flight 的 call 发CancelCall,清理 dispatcher 状态,关闭TaskIOScope,删除taskRuntime,写task_finished日志,并通过OnTaskTerminal通知订阅者。slot 释放已在WorkerCore.convergeTask内同步完成。 - 定时器:
app.Run每 500ms 调一次orch.TickTimers(),它依次处理 slot lease 过期、task 超时、Watch grace 到期和结果缓存清理,然后SyncExecutorSnapshots、CheckWedgedExecutors、execAdapter.CheckHeartbeatTimeouts。 - shadow 转发:配置
ShadowSubmitTargetAddr后,SubmitTaskShadowForwarder在主提交成功后按采样率把请求镜像到另一 Worker;接收端通过isShadowSubmit识别并剥离ResultHandler。
orchestrator_wedge.go 中的 CheckWedgedExecutors 是一个兜底:某个 executor 仍在心跳但 LastSchedulerActiveAtMs 超过 SchedulerStallTimeoutMs 没有推进且报告 RunningCallID,就调用 SetExecutorKiller 注入的钩子硬杀它,让 Pool Manager 拉起新进程。SchedulerStallTimeoutMs 由 app.Run 从 CallDeadlineMs 推导并钳到安全下限。
状态与生命周期
slot 状态机
SlotState 只有三个值(core/types.go):FREE、ALLOCATED、RUNNING。slot 数量固定为 taskSlots,SlotID 形如 slot-N。
HandleRequestTaskSlot先做幂等检查(同taskId仍活跃则返回稳定快照;已终态且在保留窗口内则返回终态),再检查freeCapacity() = TaskSlots - (ALLOCATED + RUNNING),最后线性扫描第一个FREEslot 分配。SubmitTask不带slotId时,core 内联执行一次HandleRequestTaskSlot,TTL 用DefaultTTLMs。- 若 lease 在
SubmitTask与HandleTaskActivationPrepared之间过期且 slot 被回收,core 把该 task 以SLOT_LOST收敛为FAILED。
task 可见状态与内部阶段
对外可见状态(commontypes.TaskState):ALLOCATED、RUNNING、SUCCEEDED、FAILED。内部阶段(core.TaskPhase):Preparing、Builder、Calls、Plugin、Terminal。所有内部阶段对外都是 RUNNING。
HandleTaskPhaseFinished 的推进规则:
Builder失败 →FAILED,失败码取 outcome 或BUILDER_FAILED。成功 →Calls。Calls失败 →FAILED。成功 → 总是输出StartWriterPhase。Plugin结束 → 终态;失败码取 outcome 或PLUGIN_FAILED,PluginResults写入ExecuteResult。- 乱序的阶段事件(
tc.Phase != phase)被静默丢弃。
时间语义
幂等语义:
RequestTaskSlot 与 SubmitTask 都按 taskId 在单 Worker 内幂等。同 taskId 已在 tasks 中时,SubmitTask 返回 RUNNING 快照;已终态且在保留窗口内则返回终态快照,不重新执行。
task 回收复用 WatchTasks 流做断连判定,没有独立的 owner lease,也没有 Cancel API。上表的 Watch grace 就是这条路径,WatchDisconnectGraceMs 默认 0,即默认关闭。设计评审见 blockx 仓库 docs/specs/2026-07-16-task-owner-lease-and-cancellation.md。
启动与关闭顺序
1
启动
Run 依次:校验 Profile → 加载配置(内置默认 → TuneDefaults → WORKER_CONFIG JSON → 环境变量)→ 初始化 tracing / Prometheus / chlog → 装配 IO backend、Function Code View、registry、Orchestrator、SubscriptionManager、GRPCServer → 启动 UDS server 和 executor PoolManager → RunDispatchLoop → 首次 SyncExecutorSnapshots → 500ms ticker → gRPC Serve → 等第一个健康 executor 出现(WorkerCore.SetReady)后才注册 etcd 并启动心跳。2
关闭(SIGINT / SIGTERM)
Phase 1 drain:停止注册 →
orch.SetDraining()(新的 RequestTaskSlot / 未激活的 SubmitTask 返回 Unavailable)→ etcd Deregister → WaitDrain 直到活跃 task 为 0 或 DrainTimeoutMs 到期。Phase 2 force:cancel lifecycle ctx → poolMgr.Stop() → execAdapter.Close() → CloseIOScope() → subMgr.CloseAll() → GracefulStop(5s 后 Stop)。配置
配置结构是app.WorkerFullConfig(internal/worker/app/config.go)。加载顺序:DefaultWorkerFullConfig() → Profile.TuneDefaults → WORKER_CONFIG 指向的 JSON 文件(LoadWorkerFullConfigFrom,只覆盖文件里出现的字段)→ 环境变量(app.Run 中逐项 envOverride*)。显式写的文件 / env 值总是赢,这也是回滚 profile 默认值的通道。
core.Config.Validate 要求 taskSlots > 0、taskDeadlineMs <= maxTaskTimeoutMs、defaultRetryJitterPercent 在 0 到 5 之间。dispatcher 相关字段(dispatcher.*)见 Call 执行子系统。
扩展点
- 新增 Call Builder:在
internal/plugin/callbuilder/实现cbtypes.CallBuilder(可选StreamingCallBuilder),在cbtypes中定义CallBuilderName,然后在internal/worker/app/app.go的newBuilderRegistryswitch 中加一个 case,并把名字加进需要它的cmd/*/main.goProfile 的Builders。未列入 Profile 的 builder 不会注册,task 会以BUILDER_NOT_FOUND失败。详见 Plugin 系统。 - 新增 Writer Plugin:实现
evtypes.WriterPlugin,在newPluginRegistry中加 case,并加进 Profile 的Plugins。若插件依赖特定后端地址,在validateProfileRuntimeConfig加校验。 - 新增入口 / 部署形态:只写一个新的
cmd/<name>/main.go,声明app.Profile并调用app.Run。不要在入口里手搓 wiring。Profile.Validate目前只接受block和bundle两种Deployment。 - 新增配置项:在
WorkerFullConfig加字段和默认值,在app.Run加envOverride*,必要时在Validate加约束。 - 改 slot / 阶段 / 幂等语义:只改
internal/worker/core/worker.go,并在internal/worker/core/worker*_test.go补 core unit 测试。adapter 不应出现新的状态判定。 - 改阶段执行副作用:
orchestrator_phases.go。任何触碰dispCore的新路径必须遵守orchestrator_actor.go文件头的 I1–I5 不变量,actor_invariants_test.go会拦截违规。 - 新增对外 RPC:改
api/grpc/worker/worker.proto,重新生成workerpb,在grpc_server.go加 handler,pb 转换放wire.go。
测试
按docs/specs/worker-test-organization.md 分三层:
- core unit(
internal/worker/core/):worker_test.go、worker_activation_test.go、worker_watch_test.go覆盖WorkerCore;dispatcher_*_test.go覆盖DispatcherCore。不引入 transport 或进程。 - adapter integration(
internal/worker/adapters/):orchestrator_*_test.go(lifecycle、phases、dispatch、stream、budget、deadline_repro、subcall、send_failure 等切片)、grpc_server_test.go、stream_subscriber_test.go、etcd_publisher_test.go、actor_invariants_test.go。使用 fake executor adapter,不起 Python。internal/worker/app/*_test.go覆盖 Profile、配置解析和 etcd 注册装配。 - process e2e(
cmd/worker/worker_process_*_test.go、cmd/bundle_worker/*_test.go):拉起真实 worker 进程和 Python executor,验证对外 gRPC contract(activation、watch、shutdown、timing、subcall、io、function code 等切片)。
<boundary>_test.go 或 <boundary>_<slice>_test.go,一个文件只属于一个实现边界和一个测试层级。更多命令见 测试组织与命令。
相关文档
blockx 仓库中的 spec:docs/specs/worker.md:Worker 详细设计。协议、状态模型、执行流程、Executor Pool Manager、心跳与失败语义。docs/specs/architecture.md§4.2.1 / §4.2.2:Worker 在整体架构中的职责,以及阶段级准入、task 内放量窗口、backend admission 三层并发控制。docs/specs/worker-test-organization.md:Worker 测试组织标准。docs/specs/2026-07-16-task-owner-lease-and-cancellation.md:Watch 断连回收方案(watchDisconnectGraceMs)的评审记录。docs/specs/2026-07-28-registry-prefix-migration.md:workerRegistryPrefix的格式与迁移。
Call 执行子系统
dispatcher、executor adapter、subcall 与 wedge 兜底
Task Resource Coordinator
谁调用 RequestTaskSlot,以及 Coordinator 如何消费心跳
Plugin 系统
Call Builder 与 Writer Plugin 接口
Function Code View
taskCodeEpoch 与代码快照的来源
IO 访问子系统
TaskIOScope 与 backend admission
协议与接口
gRPC / etcd 协议细节