ReserveWorkerSlot(taskId),它从 etcd 心跳重建的 Worker 视图里选一个候选,向该 Worker 调 RequestTaskSlot,然后把 workerAddr + slotId 交还给 Client。之后 Client 直接向 Worker 提交 task,Coordinator 不再参与。它在整体架构中的位置见 架构总览,主链路见 Task 生命周期。
职责与边界
负责:- watch etcd 中的 Worker 注册与心跳,维护一份本地、可丢弃的
WorkerView。 - 对每次
ReserveWorkerSlot做校验、就绪判定、候选过滤与排序,然后串行尝试少量候选。 - 把 Worker 返回的
slotId原样转交 Client;不解释slotId的内部含义。 - 基于最近失败对 Worker 做本地惩罚和短暂熔断。
- 解析 task payload、生成
taskId、接收SubmitTask或TaskResult。 - 维护中心化 slot 账本、pending 队列或全局
taskId去重。 - 向 etcd 写任何数据;etcd 只承载 Worker 发现与心跳。
- Worker 本地 slot 状态机是唯一真相;Coordinator 重启只丢失视图缓存和路由偏置,不影响正确性。
CoordinatorCore是 Sans-IO 的纯决策引擎:不读时钟、不做网络,所有方法显式接收now(UTC 毫秒)。- 首轮全量 Worker 视图建立前,以及 watch 断流超过宽限期后,
ReserveWorkerSlot一律返回Unavailable。 - 单次请求串行尝试候选,不做 hedged request;结果不确定时只对同一 Worker、同一
taskId重试,绝不切换 Worker。 - 同一套代码用不同
app.Profile起两套进程:cmd/coordinator(block 集群)和cmd/bundle_coordinator(bundle 集群)。bundle 集群的细节见 Bundle 集群。
代码位置
核心类型与接口
core
core.Config(internal/coordinator/core/coordinator.go):选址与保护参数,DefaultConfig()给默认值,Validate()要求WatchGracePeriodMs <= HeartbeatTimeoutMs。core.CoordinatorCore:持有map[string]*WorkerView、initialized、watchDisconnectedAtMs。事件入口:ApplyFullWorkerList、ApplyWorkerUpdate、ApplyWorkerRemoval、MarkWatchDisconnected、MarkWatchReconnected。决策出口:IsReady、SelectCandidate。反馈入口:RecordSlotSuccess、RecordSlotRejection、RecordSlotBackpressure、RecordSlotTimeout。core.WorkerView(internal/coordinator/core/worker_view.go):内嵌etcd.WorkerHeartbeat,加三个只在本地存在、不写回 etcd 的字段。
WorkerView.UsedTotal()=ReservedSlots + RunningTasks;WorkerView.FreeCapacity()=TaskSlots - UsedTotal()。core.CandidateResult{WorkerAddr, Err}与core.CoordinatorError{Code, Message}:SelectCandidate的输出,Code是commontypes.ErrorCode。core.ValidateReserveRequest(taskID)(internal/coordinator/core/validate.go):只检查taskId非空。
adapters
adapters.CoordinatorServer(grpc_server.go):实现coordinatorpb.CoordinatorServiceServer,持有core、SlotReserver和一把sync.RWMutex;Lock/Unlock暴露给 watcher 用于写 core。adapters.HandlerConfig{MaxTimeoutRetries}:不确定结果时对同一 Worker 的重试次数,默认1。adapters.SlotReserver接口:RequestTaskSlot(ctx, workerAddr, req) (*commontypes.TaskSnapshot, *SlotError),测试里用 mock 替换。adapters.WorkerRPCClient(worker_rpc.go):SlotReserver的 gRPC 实现,按workerAddr缓存一条grpc.ClientConn;grpcErrToSlotError决定SlotError.MayHaveAllocated。adapters.SlotError{Code, Message, MayHaveAllocated}:MayHaveAllocated=true表示请求可能已到达 Worker,必须留在同一 Worker 重试。adapters.EtcdWorkerWatcher/EtcdWatcherConfig{Endpoints, KeyPrefixes, DialTimeout, RetryDelay}(etcd_watcher.go):Run(ctx)阻塞运行 list+watch。adapters.BundleCoordinatorServer(bundle_grpc_server.go):bundle proto 到CoordinatorServer的薄转换层。
app 与 api
app.Profile{Deployment, Service, WorkerRegistryPrefixes}(internal/coordinator/app/app.go):入口身份。Deployment只能是"block"或"bundle",并决定注册哪个 gRPC service、允许哪种 registry kind。app.Run(p Profile):装配 obs、core、WorkerRPCClient、CoordinatorServer、etcd client、watcher、gRPC server,阻塞到收到 SIGINT/SIGTERM。app.CoordinatorRuntimeConfig(config.go):进程级配置聚合,LoadCoordinatorRuntimeConfigFromEnv(p)从环境变量覆盖并Validate()。etcd.WorkerHeartbeat(api/etcd/types.go):Worker 写入 etcd 的 JSON 值,字段包括workerAddr、taskSlots、reservedSlots、runningTasks、cpuPercent、memoryPercent、memoryAvailableBytes、memoryTotalBytes、processRssBytes、errorRate、lastHeartbeatMs。etcd.ParseRegistryPrefix与常量LegacyBlockWorkerRegistryPrefix(/blockx/workers/)、LegacyBundleWorkerRegistryPrefix(/blockx/bundle-workers/)、ProdDefaultBlockWorkerRegistryPrefix、ProdDefaultBundleWorkerRegistryPrefix、TestDefaultBundleWorkerRegistryPrefix(api/etcd/registry.go)。
数据流 / 执行流程
etcd key 是<prefix> + workerAddr,value 是 WorkerHeartbeat JSON。Worker 侧用带 lease 的 Put 定期刷新,写入逻辑在 internal/worker/adapters/etcd_publisher.go(lease TTL 默认 10s,心跳间隔见 Worker 配置 HeartbeatIntervalMs)。Coordinator 只读。
CoordinatorServer.ReserveWorkerSlot 的编排逻辑按以下顺序执行(internal/coordinator/adapters/grpc_server.go):
1
校验与就绪判定
core.ValidateReserveRequest 拒绝空 taskId;core.IsReady(now) 为 false 时直接返回 Unavailable,此时未向任何 Worker 发请求。2
选候选
在读锁下调用
core.SelectCandidate。硬过滤依次为:已在 excludeAddrs、心跳超时(now - LastHeartbeatMs > HeartbeatTimeoutMs)、CPUPercent/MemoryPercent/ErrorRate 超红线、熔断未到期、FreeCapacity() <= 0。剩下的按 candidateRanksBefore 排序:FreeCapacity 大者优先,再比 LastHeartbeatMs 新者、Penalty 小者、最后按 WorkerAddr 字典序。3
向 Worker 预占
WorkerRPCClient.RequestTaskSlot 带 WorkerRPCTimeout 调 Worker。grpcErrToSlotError 把 gRPC 状态码分成两类:InvalidArgument/ResourceExhausted/FailedPrecondition/NotFound 是明确拒绝;DeadlineExceeded/Canceled/Aborted/Internal/Unknown 以及超时型 Unavailable 标记 MayHaveAllocated=true。返回成功但 SlotID 为空同样按不确定处理。4
反馈与收敛
成功调用
RecordSlotSuccess 清零失败计数。明确拒绝时:Worker 返回 ResourceExhausted(映射为 ErrNoSlot)走 RecordSlotBackpressure,不计入熔断;其他拒绝走 RecordSlotRejection。不确定结果先在同一 Worker 重试,仍不确定则 RecordSlotTimeout 并返回 Aborted;重试中拿到明确拒绝则换下一个候选。CoordinatorCore 不做锁保护;CoordinatorServer.mu 是唯一的同步点,watcher 通过 srv.Lock()/Unlock() 写 core,handler 用读锁选候选、写锁记录反馈。
对外错误码
对外协议是 gRPC。业务错误码到 gRPC 状态码的映射在coordCodeToGRPC:
重试后仍无法确认时返回 gRPC
Aborted;coordCodeToGRPC 把 ErrSlotUncertain 映射为 Aborted,测试工具 internal/testutil/rpc_types.go 反向解释为 ErrSlotUncertain。调用方应把 Aborted 当作”可能已预占,至少等一个 slot TTL 再重试”,把 Unavailable 当作”可立即重试”。状态与生命周期
CoordinatorCore 只有两个全局状态位,没有显式状态机:
initialized:ApplyFullWorkerList第一次调用后为 true,之后不再回退。watcher 只有在所有配置 prefix 都完成初始 list 后才调用它。watchDisconnectedAtMs:任一 prefix 的 watch 断开时由MarkWatchDisconnected(now)记录首次断开时间;IsReady与SelectCandidate在now - watchDisconnectedAtMs > WatchGracePeriodMs时返回Unavailable。所有 prefix 重新 list 成功且都处于连接态后,MarkWatchReconnected清零。
RecentFailureCount与Penalty:RecordSlotRejection/RecordSlotTimeout各加一,Penalty = float64(RecentFailureCount)。CircuitBreakerDeadlineMs:RecentFailureCount >= CircuitBreakerThreshold(默认 3)时设为now + CircuitBreakerCooldownMs(默认 10s),期间该 Worker 被硬过滤。RecordSlotSuccess把三者一起清零;RecordSlotBackpressure对三者都不动。ApplyFullWorkerList会保留幸存 Worker 的这三个字段;ApplyWorkerRemoval删除后重新出现的 Worker 从零开始。
- slot TTL(
SlotTTLMs,默认 5s)是 Coordinator 在RequestTaskSlot里传给 Worker 的预占窗口,不是 task 执行超时;Coordinator 不再解释它。 - 单次
ReserveWorkerSlot没有独立的总时间预算;上界由MaxCandidates × (1 + MaxTimeoutRetries) × WorkerRPCTimeout和调用方 ctx 共同决定,默认最坏约3 × 2 × 3s。 - 心跳新鲜度用 Coordinator 本地时钟与 Worker 写入的
lastHeartbeatMs相减,机器间时钟偏差会直接影响过滤结果。
配置
所有配置来自internal/coordinator/app/config.go,默认值取自 DefaultCoordinatorRuntimeConfig、core.DefaultConfig、adapters.DefaultHandlerConfig 与 adapters.DefaultEtcdWatcherConfig。
可观测性相关的
BLOCKX_PROMETHEUS_LISTEN_ADDR、BLOCKX_PROMETHEUS_PATH、CHAINTABLE_LOG_BROKERS、CHAINTABLE_LOG_TOPIC 由 internal/obs 统一处理,见 可观测性。COORDINATOR_CPU_PROFILE / COORDINATOR_TRACE 让 app.Run 在整个进程生命周期写 pprof / runtime trace 文件,只用于性能排查。
Profile 默认 registry:
cmd/coordinator:/blockx/workers/。cmd/bundle_coordinator:/blockx/bundle-workers/与/blockx/prod/lanes/default/bundle-workers/。- 测试集群用
COORDINATOR_REGISTRY_PREFIXES=/blockx/test/lanes/default/bundle-workers/整体替换。
EtcdWorkerWatcher 按 WorkerAddr 合并:LastHeartbeatMs 更新者胜出,相同则 v2 格式胜过 legacy,再相同按 prefix 字符串比较;一个 Worker 只有在所有 prefix 下都消失后才从视图删除。详见 blockx 仓库 docs/specs/2026-07-28-registry-prefix-migration.md。
扩展点
- 改选址策略:只动
internal/coordinator/core/coordinator.go的SelectCandidate与candidateRanksBefore,并在coordinator_test.go补TestSelectCandidate_*。core 没有 IO,可以纯函数式测试。当前是确定性排序;心跳里的绝对内存字段(memoryAvailableBytes)只被解析,不参与过滤和排序。 - 新增一种 Worker 反馈:在 core 加
RecordSlotXxx,在adapters/grpc_server.go的recordSlotOutcome里按SlotError.Code分派。RecordSlotBackpressure就是这样加进来的。 - 改 Worker 错误分类:
adapters/worker_rpc.go的grpcErrToSlotError与 Worker 侧workerCodeToGRPC(internal/worker/adapters/wire.go)必须保持互逆;改一边就要看另一边(见 Worker)。 - 新增部署形态:加一个
cmd/<name>/main.go声明app.Profile,在app.registerCoordinatorService与validateRegistryPrefixes里加分支,并在api/etcd/registry.go的ParseRegistryPrefix里允许新的 worker kind。 - 换 Worker RPC 实现:实现
adapters.SlotReserver接口即可,CoordinatorServer不依赖具体 client。当前WorkerRPCClient对每个workerAddr只建一条grpc.ClientConn,不使用internal/sdk/grpcclient的多连接 balancer:指向单个 Worker 的点对点 client 不启用多连接策略。背景见 blockx 仓库docs/specs/2026-07-23-grpc-client-balancing.md。 - 新增配置项:在
core.Config/adapters.HandlerConfig/adapters.EtcdWatcherConfig加字段和Validate规则,再在app/config.go的LoadCoordinatorRuntimeConfigFromEnv加环境变量解析,并补config_test.go。
测试
internal/coordinator/core/coordinator_test.go:SelectCandidate每个过滤条件与排序规则、熔断触发/恢复、IsReady的宽限期语义、ApplyFullWorkerList保留本地字段。internal/coordinator/adapters/handler_test.go:用 mockSlotReserver覆盖编排循环(超时后换候选、重试仍超时返回Aborted、全部拒绝返回ResourceExhausted、空SlotID重试),以及真实 gRPC server 的错误码与压缩行为。internal/coordinator/adapters/worker_rpc_test.go:grpcErrToSlotError的MayHaveAllocated分类(连接拒绝 vs 超时/取消/Aborted)。internal/coordinator/adapters/etcd_watcher_test.go:内嵌 etcd 上的 list+watch、多 prefix 合并去重、断连/重连的就绪状态、不可解析 value 跳过。internal/coordinator/app/config_test.go、profile_test.go:环境变量解析与Profile校验。e2e/contract/coordinator-etcd/:etcd 晚于 Coordinator 启动时先Unavailable后恢复;watch 在宽限期内恢复仍可路由;超过宽限期变Unavailable。e2e/contract/coordinator-worker/:明确传输失败切换下一个 Worker、RequestTaskSlot超时返回SlotUncertain(Aborted)、slot 过期后SubmitTask被拒、Worker 重启后旧slotId失效、统一 slot 池满时拒绝。
internal/testutil 的 StartCoordinator / StartBundleCoordinator 启动真实进程,用 COORDINATOR_* 环境变量缩短时序参数。整体测试组织见 测试。
相关文档
blockx 仓库中的 spec:docs/specs/task-resource-coordinator.md:Coordinator 详细设计,协议、错误语义、选址规则、一致性边界。docs/specs/architecture.md§4.1:控制面在整体架构中的定位与三条边界。docs/specs/2026-07-28-registry-prefix-migration.md:legacy / v2 registry prefix、多 prefix 合并与发布顺序。docs/specs/2026-07-23-bundle-coordinator-api.md:bundle Coordinator 独立 service 身份的由来。docs/specs/2026-07-23-grpc-client-balancing.md:多连接 balancer 的适用范围。docs/specs/worker.md:RequestTaskSlot与 slot 状态机的权威定义。