Skip to main content
Worker 是 BlockX 的执行面节点。它在本机维护固定大小的 task slot 池,接收 Coordinator 或直连 Client 提交的 task,按 Builder → Calls → Plugin 的顺序执行,并在短期窗口内保留结果供查询和订阅。它在整体架构中的位置见 架构总览,task 端到端链路见 Task 生命周期

职责与边界

Worker 负责:
  • 对外提供 gRPC WorkerServiceapi/grpc/worker/worker.proto):RequestTaskSlotSubmitTaskGetTaskResultWatchTasksWatchTasks 是 server-streaming RPC,Worker 只监听一个 TCP 端口。协议细节见 协议与接口
  • 维护本地 slot 表和 task 索引,做原子准入和 taskId 级幂等。
  • 为每个 task 固定一个 taskCodeEpoch,串行跑 Call Builder,把 call list 交给 dispatcher 并行执行,最后串行跑 Writer Plugin。
  • 收敛 TaskResult,推送给 WatchTasks 订阅者,并按 resultRetentionMs 保留。
  • 向 etcd 注册并周期上报心跳,供 Coordinator 选址。
Worker 不负责:
  • 生成 taskId、构造 payload、DAG 编排或定时触发(这些在上游)。
  • 跨 Worker 的去重、全局结果查询或崩溃后的本地恢复。重启后旧 slot、旧结果和订阅全部丢失。
  • call 级调度细节(dispatcher、executor adapter)。这部分见 Call 执行子系统
核心不变量:
  • WorkerCore 是纯内存、Sans-IO 的决策引擎。它不读时钟,每个方法都显式接收 now
  • slot 表是容量和 task 状态迁移的唯一真相;etcd 心跳只是它的派生视图。
  • slotTabletaskIndex 在同一次方法调用里一起更新,不存在撕裂状态。
  • 一旦 task 成功激活,后续所有失败(Builder、call、Plugin、超时、Watch 断连)都收敛为 TaskResult 里的终态,不再返回提交层错误。
  • adapter 不维护第二份 slot 或 phase 真相。它只执行 core 输出的命令,再把结果作为事件回送。

代码位置

核心类型与接口

core(internal/worker/core/
  • WorkerCoreworker.go):单线程状态机。持有 slots []SlotEntrytaskIndex map[string]*TaskIndexEntrytasks map[string]*TaskContextreadydraining
  • 输入事件(都在 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):PrepareTaskActivationStartBuilderPhaseStartCallPhase{CallList, Streaming}StartWriterPhase{Outputs}ConvergeTask{State, Result}。它们嵌在结果类型 RequestTaskSlotResultSubmitTaskResultActivationResultPhaseTransitionResultLeaseExpiredResultGetTaskResultOutput 里返回。
  • TaskContexttypes.go):TaskIDSlotIDTaskCodeEpochPhaseDeadlineMsWatchAttached / WatchEverAttached / WatchDisconnectDeadlineMsCallFailurePolicyCallRetryPolicyAcceptedAtMsAdapterTaskCtx,以及三个阶段的 PhaseOutcome
  • DeriveHeartbeat(now) etcd.WorkerHeartbeat:从 slot 表派生 ReservedSlots / RunningTasks
adapters(internal/worker/adapters/
  • GRPCServergrpc_server.go):实现 RequestTaskSlot / SubmitTask / GetTaskResult,把 WatchTasks 委托给 SubscriptionManagerworkerErrToStatuswire.go)把 core.WorkerError 映射为 gRPC status,业务错误码经 workerCodeToGRPC 映射:InvalidArgument → InvalidArgumentNoSlot → ResourceExhaustedUnavailable → UnavailableInvalidSlot → FailedPreconditionNotFound → NotFoundSlotUncertain → Aborted
  • Orchestratororchestrator.go):桥接 WorkerCoreDispatcherCore。入口方法 RequestTaskSlotSubmitTaskGetTaskResultTickTimersDeriveHeartbeatSetDrainingWaitDrainRunDispatchLoopSyncExecutorSnapshotsCheckWedgedExecutors
  • OrchestratorDepsorchestrator.go):构造依赖,包括 BuilderRegistryPluginRegistryIOScopeAdmissionFunctionCodeViewProviderEpochResolverExecAdapterAuditGateStreamBuild
  • taskRuntimetask_runtime.go):adapter 侧的 task 运行时记录,持有 ioTaskCtxTaskIOScopectx / cancelfunctionCodeView、流式生产句柄和 TaskStatistic。它不维护 phase 或 slot 真相。
  • SubscriptionManagerstream_subscriber.go):每条 WatchTasks 流一个 streamSubscriber;把首条 attach / 最后一条 detach 翻译为 Orchestrator.HandleWatchAttached / HandleWatchDetached
  • EtcdPublisheretcd_publisher.go):RegisterRunHeartbeatLoopDeregister。key 为 KeyPrefix + workerAddr,lease TTL 默认 10s。
  • ExecutorAdapter 接口(exec_adapter.go):Orchestrator 依赖的 UDS executor 面,生产实现是 executor.Adapter
app(internal/worker/app/
  • Profileapp.go):Deploymentblock / bundle,只做观测标识)、Service(tracing / Prometheus 服务名)、WorkerRegistryPrefix(etcd 注册前缀)、UsageServiceBuildersPluginsTuneDefaults(只改内置默认值)。
  • Run(p Profile):唯一的装配和启停序。

数据流 / 执行流程

Worker 采用”输入事件 → core 决策 → adapter 执行命令 → 事件回送”的模式。Orchestratoro.mu 保护 WorkerCore,用 runtimeMu 保护 runtimesDispatcherCore 只在 dispatch actor(RunDispatchLoop goroutine)上被修改,其他 goroutine 通过 postDispatchEvent 投递事件。锁序和 actor 不变量写在 orchestrator_actor.go 文件头注释,并由 actor_invariants_test.go 强制。 几个关键点:
  • SubmitTask 在返回前只做到 HandleSubmitTask。core 立即创建 TaskContextPhase = Preparing),使重复提交命中幂等路径;激活和阶段执行在 activateAndRun goroutine 中异步进行。
  • 激活失败(epoch 解析失败、Pin 失败、NewTaskScope 失败)走 HandleTaskActivationFailed,失败码为 ACTIVATION_FAILEDIO_SCOPE_FAILED
  • Builder 阶段:runBuilderPhase 先取 builderSem,再从 BuilderRegistryFunctionCallConfig.Type。查不到直接以 BUILDER_NOT_FOUND 收敛。若 builder 实现了 StreamingCallBuilderStreamBuild.Enabled,Builder 阶段只跑 PrepareStream,回送 CallStreamMarker,core 输出 StartCallPhase{Streaming: true},随后 runCallStreamPhasescanSem 下启动生产 goroutine。
  • Calls 阶段:runCallPhaseexecutorSem,在锁外用 dispCore.PrepareCallPhaseWithDigests 准备 call 上下文,再投递 evCallPhaseStarted 给 actor。call 级调度、重试、subcall 归 Call 执行子系统。actor 收到 ConvergeTaskCalls 命令时调用 HandleTaskPhaseFinished(PhaseCalls)
  • Plugin 阶段:runWriterPhasepluginSem,按 ResultHandler.TypePluginRegistry。未配置 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 到期和结果缓存清理,然后 SyncExecutorSnapshotsCheckWedgedExecutorsexecAdapter.CheckHeartbeatTimeouts
  • shadow 转发:配置 ShadowSubmitTargetAddr 后,SubmitTaskShadowForwarder 在主提交成功后按采样率把请求镜像到另一 Worker;接收端通过 isShadowSubmit 识别并剥离 ResultHandler
orchestrator_wedge.go 中的 CheckWedgedExecutors 是一个兜底:某个 executor 仍在心跳但 LastSchedulerActiveAtMs 超过 SchedulerStallTimeoutMs 没有推进且报告 RunningCallID,就调用 SetExecutorKiller 注入的钩子硬杀它,让 Pool Manager 拉起新进程。SchedulerStallTimeoutMsapp.RunCallDeadlineMs 推导并钳到安全下限。

状态与生命周期

slot 状态机

SlotState 只有三个值(core/types.go):FREEALLOCATEDRUNNING。slot 数量固定为 taskSlotsSlotID 形如 slot-N
  • HandleRequestTaskSlot 先做幂等检查(同 taskId 仍活跃则返回稳定快照;已终态且在保留窗口内则返回终态),再检查 freeCapacity() = TaskSlots - (ALLOCATED + RUNNING),最后线性扫描第一个 FREE slot 分配。
  • SubmitTask 不带 slotId 时,core 内联执行一次 HandleRequestTaskSlot,TTL 用 DefaultTTLMs
  • 若 lease 在 SubmitTaskHandleTaskActivationPrepared 之间过期且 slot 被回收,core 把该 task 以 SLOT_LOST 收敛为 FAILED

task 可见状态与内部阶段

对外可见状态(commontypes.TaskState):ALLOCATEDRUNNINGSUCCEEDEDFAILED。内部阶段(core.TaskPhase):PreparingBuilderCallsPluginTerminal。所有内部阶段对外都是 RUNNING HandleTaskPhaseFinished 的推进规则:
  • Builder 失败 → FAILED,失败码取 outcome 或 BUILDER_FAILED。成功 → Calls
  • Calls 失败 → FAILED。成功 → 总是输出 StartWriterPhase
  • Plugin 结束 → 终态;失败码取 outcome 或 PLUGIN_FAILEDPluginResults 写入 ExecuteResult
  • 乱序的阶段事件(tc.Phase != phase)被静默丢弃。

时间语义

幂等语义:RequestTaskSlotSubmitTask 都按 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 → 加载配置(内置默认 → TuneDefaultsWORKER_CONFIG JSON → 环境变量)→ 初始化 tracing / Prometheus / chlog → 装配 IO backend、Function Code View、registry、OrchestratorSubscriptionManagerGRPCServer → 启动 UDS server 和 executor PoolManagerRunDispatchLoop → 首次 SyncExecutorSnapshots → 500ms ticker → gRPC Serve → 等第一个健康 executor 出现(WorkerCore.SetReady)后才注册 etcd 并启动心跳。
2

关闭(SIGINT / SIGTERM)

Phase 1 drain:停止注册 → orch.SetDraining()(新的 RequestTaskSlot / 未激活的 SubmitTask 返回 Unavailable)→ etcd DeregisterWaitDrain 直到活跃 task 为 0 或 DrainTimeoutMs 到期。Phase 2 force:cancel lifecycle ctx → poolMgr.Stop()execAdapter.Close()CloseIOScope()subMgr.CloseAll()GracefulStop(5s 后 Stop)。

配置

配置结构是 app.WorkerFullConfiginternal/worker/app/config.go)。加载顺序:DefaultWorkerFullConfig()Profile.TuneDefaultsWORKER_CONFIG 指向的 JSON 文件(LoadWorkerFullConfigFrom,只覆盖文件里出现的字段)→ 环境变量(app.Run 中逐项 envOverride*)。显式写的文件 / env 值总是赢,这也是回滚 profile 默认值的通道。 core.Config.Validate 要求 taskSlots > 0taskDeadlineMs <= maxTaskTimeoutMsdefaultRetryJitterPercent 在 0 到 5 之间。dispatcher 相关字段(dispatcher.*)见 Call 执行子系统

扩展点

  • 新增 Call Builder:在 internal/plugin/callbuilder/ 实现 cbtypes.CallBuilder(可选 StreamingCallBuilder),在 cbtypes 中定义 CallBuilderName,然后在 internal/worker/app/app.gonewBuilderRegistry switch 中加一个 case,并把名字加进需要它的 cmd/*/main.go Profile 的 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 目前只接受 blockbundle 两种 Deployment
  • 新增配置项:在 WorkerFullConfig 加字段和默认值,在 app.RunenvOverride*,必要时在 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 unitinternal/worker/core/):worker_test.goworker_activation_test.goworker_watch_test.go 覆盖 WorkerCoredispatcher_*_test.go 覆盖 DispatcherCore。不引入 transport 或进程。
  • adapter integrationinternal/worker/adapters/):orchestrator_*_test.go(lifecycle、phases、dispatch、stream、budget、deadline_repro、subcall、send_failure 等切片)、grpc_server_test.gostream_subscriber_test.goetcd_publisher_test.goactor_invariants_test.go。使用 fake executor adapter,不起 Python。internal/worker/app/*_test.go 覆盖 Profile、配置解析和 etcd 注册装配。
  • process e2ecmd/worker/worker_process_*_test.gocmd/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.mdworkerRegistryPrefix 的格式与迁移。
站内页面:

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 协议细节