Compare commits

...

3 Commits

Author SHA1 Message Date
youfak
081e172970 docs: plan worker control plane session runtime
Some checks are pending
ci / test (ubuntu-latest) (push) Waiting to run
ci / test (windows-latest) (push) Waiting to run
ci / race (push) Waiting to run
ci / integration (push) Waiting to run
2026-07-31 10:19:30 +08:00
youfak
a92deb9def docs: define negative worker snapshot ack 2026-07-31 07:02:50 +08:00
youfak
3e6d459777 docs: design worker session control plane 2026-07-31 06:45:53 +08:00
2 changed files with 1591 additions and 0 deletions

File diff suppressed because it is too large Load Diff

View File

@ -0,0 +1,310 @@
# Worker 控制面会话与运行态接收设计
## 1. 目标
实现 `WorkerControlPlane` 的第一个端到端交付单元,使 Controller 能通过 gRPC
1. 注册 Gateway Worker 并签发可过期、可替换的 Session。
2. 校验 Worker 对已发布 Snapshot 的正向或负向 ACK。
3. 接收完整稀疏 Runtime 报告,并把空报告作为心跳。
4. 在多个 Controller 副本之间通过 Redis 共享 Session、Snapshot 元数据和报告
序号栅栏。
本次实现 `RegisterWorker`、`AcknowledgeSnapshot` 和 `ReportRuntime` 三个 RPC。
`WatchSnapshots``ReportOutcomes` 保留已发布的 protobuf 契约,但明确返回
`Unimplemented`,后续交付单元再实现 Snapshot payload/stream、Gateway 客户端和
Outcome 去重。
## 2. 当前基础与约束
仓库已经具备以下基础:
- `controlplane.proto` 已定义 Worker 注册、Snapshot、ACK、Outcome 和 Runtime。
- `workerruntime` 已定义 Session、完整稀疏报告和内存参考实现。
- `redisactivity` 已原子维护 Worker Session、报告序号、报告 TTL 和 ownership
fence。
- Gateway Snapshot Store 已能生成与当前本地快照一致的 Runtime 报告。
- Kubernetes 清单已预留 Controller `8443` 控制面端口。
现有 `SessionWriter.ReplaceSession` 只能保存已经 ACK 的 Session不能表达注册后、
首次 ACK 前的状态;现有 `ErrStaleReport` 同时表示报告序号落后和 Snapshot 不匹配,
不足以给 Worker 返回准确恢复指令。本次必须修正这两个领域边界,而不是只在 gRPC
Handler 中增加特殊分支。
必须保持以下不变量:
- Gateway 请求热路径不访问 Redis、PostgreSQL 或 gRPC。
- Session、Snapshot 元数据和 Runtime 只进入 Redis 或节点内存,不进入 PostgreSQL。
- Worker 自报的 Session、版本、epoch 和 checksum 不能直接成为权威状态。
- Session ID、证书身份、Proxy 地址和 Client 信息不得成为 Prometheus 标签。
## 3. 组件边界
### 3.1 生成契约
`api/proto/controlplane/v1/controlplane.proto` 生成并提交:
- `gen/controlplane/v1/controlplane.pb.go`
- `gen/controlplane/v1/controlplane_grpc.pb.go`
`go.mod` 直接固定 `google.golang.org/grpc``google.golang.org/protobuf`,并用 Go
tool 依赖固定 `protoc-gen-go``protoc-gen-go-grpc`。生成脚本必须可重复执行CI
验证 descriptor 编译和已提交生成文件没有漂移。
### 3.2 公共应用服务
新增 `internal/controller/worker`,向传输层公开以下应用接口:
```go
type Service interface {
Register(context.Context, RegisterCommand) (Registration, error)
Acknowledge(context.Context, SnapshotAcknowledgement) error
ReportRuntime(context.Context, workerruntime.Report) (RuntimeDecision, error)
}
```
命令和结果不引用 protobuf 类型。gRPC DTO 映射、状态码映射和证书身份校验留在
传输适配层,领域服务可以在不启动网络的情况下测试。
服务依赖以下公共端口:
```go
type ControlStore interface {
CurrentOwnershipEpoch(context.Context) (uint64, error)
OpenSession(context.Context, workerruntime.Session, time.Duration) error
RecordIssuedSnapshot(context.Context, workerruntime.SnapshotReference, time.Duration) error
AcknowledgeSnapshot(context.Context, workerruntime.SnapshotAcknowledgement, time.Duration) error
ReplaceRuntime(context.Context, workerruntime.Report, time.Duration) error
}
```
`RecordIssuedSnapshot` 本轮用于契约和集成测试,并作为下一轮 Snapshot Publisher 的
稳定接入点。Controller 启动不会伪造已发布 Snapshot。
### 3.3 Redis 权威存储
`redisactivity.Adapter` 实现同一 `ControlStore`。新增的 Session 与 Snapshot 元数据
继续使用 `{activity}` hash tag与当前 ownership 和 runtime key 位于同一 Redis
Cluster slot。
每个 Worker 只保留当前已发布 SnapshotReference
- `worker_id`
- `version`
- `ownership_epoch`
- `checksum`
- `expires_at`
引用按 `(ownership_epoch, version)` 严格单调前进。相同 tuple 只允许使用
相同 checksum 幂等续期;同 tuple 更换 checksum 或倒退 tuple 均原子失败。
发布新引用会覆盖旧引用。旧版本 ACK 必须失败Worker 后续通过完整 Snapshot 重同步。
Snapshot payload 不在本轮写入 Redis下一轮会评估 payload 存储和流式分发方式。
全局 ownership epoch 是 namespace 内的持久单调标量,不跟随 ownership TTL
过期。既有 namespace 在首次控制面读取或下次 Assign/Renew 时移除该 key
的历史 TTL避免 epoch 回退产生 ABA。
## 4. Session 状态机
Session 有以下状态:
```text
[不存在/过期]
|
| RegisterWorker
v
[已注册ACK=0]
|
| 有效正向 ACK
v
[已 ACK可接收 Runtime]
|
| 空或非空 Runtime
+--------------------> [续期]
|
| 对新 Snapshot 的负向 ACK
v
[保留旧 ACKRuntime 已禁用]
|
| 后续有效正向 ACK
+--------------------> [已 ACK可接收 Runtime]
|
| 新 instance 注册 / TTL 过期
v
[失效]
```
规则如下:
1. 每次成功注册都生成新的 128-bit 加密随机 Session ID。
2. 同一 `worker_id` 再注册会替换旧 Session并删除旧 Runtime 报告。
3. 注册后的 `acked_snapshot_version`、`acked_ownership_epoch` 和
`acked_checksum` 均为零值,`runtime_enabled=false`。
4. Runtime 只有在 Session 已完成正向 ACK 后才可接受。
5. 空 Runtime 报告表示完整稀疏计数全部为零,同时承担心跳并续期 Session。
6. Session 到期后,旧实例的 ACK 和 Runtime 均失败;使用该报告计算容量时保持
fail-closed。
`supported_protocol_version` 本轮只接受 `1`。`worker_id` 与 `instance_id` 使用
`[A-Za-z0-9][A-Za-z0-9._-]{0,127}``zone` 为空或使用同一字符集且最长 128 字节。
标签最多 32 个,键最长 64 字节、值最长 256 字节,规范化后的键值总长度最多
4 KiB禁止空键和首尾空白。`labels` 在 protobuf 契约中是 map重复 wire key
在进入 Service 前已按 protobuf map 语义归一化为最后一值Service 只验证
解码后的唯一键集。
## 5. RPC 行为
### 5.1 RegisterWorker
处理顺序:
1. 校验消息大小、协议版本、标识与标签边界。
2. 校验 mTLS Worker 身份与请求中的 `worker_id` 一致。
3. 从 Redis 读取或原子初始化全局 ownership epoch最小值为 `1`
4. 生成新 Session ID并以未 ACK 状态替换当前 Worker Session。
5. 返回 Session、ownership epoch、heartbeat interval 和 max stale age。
同一实例重复注册也签发新 Session保证网络重试不会让两个并行连接共享写权限。
### 5.2 AcknowledgeSnapshot
ACK 操作在 Redis Lua 的一个原子边界内完成:
1. Session 必须存在、未过期,并匹配 worker/session。
2. 先按 tuple 与 Session 已 ACK 上界比较,再与当前已发布
SnapshotReference 比较。落后 tuple 是 stale ACK超前 tuple 或同 tuple
不同 checksum 是 Snapshot mismatch。
3. `applied=true`ACK 只能严格前进或完全相同地幂等重放。严格前进
时清除旧 Runtime、保存新上界并设置 `runtime_enabled=true`;相同 ACK
且 Runtime 已开启时只续期 Session不清除 Runtime 或 sequence/digest fence。
若相同 ACK 上界因先前负向 ACK 处于禁用状态,则清理残留 Runtime
并重新开启。
4. `applied=false` 是已接受的负向 ACK续期当前 Session但不推进 ACK 上界、
不保留旧 Runtime设置 `runtime_enabled=false`,也不持久化自由文本错误;
下一轮 Snapshot Stream 据此继续发送完整 Snapshot。只有后续正向 ACK
才把 `runtime_enabled` 重新设为 true因此负向 ACK 后延迟到达的旧 Runtime
不能重建已清除的报告与 sequence fence。
5. 旧版本、错误 checksum、超前版本或被替换 Session 均不修改状态。
`error_code` 只接受有界枚举风格字符串;`error_message` 只做长度校验,不写日志或
持久化,避免敏感信息进入控制面状态。
### 5.3 ReportRuntime
处理顺序:
1. 校验 Session、报告序号、时间戳、计数非负、Proxy ID 唯一和批次上限。
2. Session 必须 `runtime_enabled=true`,且 Snapshot version 与 ownership epoch 必须
精确等于 Session 已 ACK 上界。
3. 相同 sequence 和相同规范化内容是幂等重放,但仍必须续期 Session
与报告 TTL较小 sequence 拒绝;相同 sequence 不同内容判定冲突。
4. 成功报告完整替换该 Worker 的稀疏 Runtime并续期 Session 与报告 TTL。
5. 返回 Controller 当前 ownership epoch本轮 `revoke_proxy_ids` 为空。
需要把错误语义拆开:
- `ErrSnapshotMismatch`:报告与 ACK 上界不一致gRPC 返回成功响应且
`require_full_snapshot=true`,不写入报告。
- `ErrStaleReport`sequence 倒退,映射 `Aborted`
- `ErrConflictingReport`:相同 sequence 内容不同,映射 `AlreadyExists`
- `ErrStaleSession`Session 被替换或过期,映射 `FailedPrecondition`
- Redis 或内部依赖错误映射 `Unavailable`,不得泄露后端错误文本。
## 6. gRPC Listener 与安全
配置新增独立的 `controlPlane`
```yaml
controlPlane:
enabled: false
listen: 127.0.0.1:8443
protocolVersion: 1
heartbeatInterval: 10s
sessionTTL: 35s
maxStaleAge: 30s
maxMessageBytes: 4194304
maxRuntimeCounters: 100000
tls:
mode: disabled
```
验证规则:
- 启用时必须有合法监听地址、正数限制和 `sessionTTL >= 3 * heartbeatInterval`
- `maxStaleAge` 不小于 heartbeat interval。
- `tls.mode` 只允许 `disabled``mtls`
- `disabled` 只允许回环地址;非回环监听必须使用 `mtls`
- `mtls` 必须提供 server cert、server key、client CA、trust domain 和 environment。
- trust domain 必须是无 port 的小写 DNS 名environment 必须是单个有界 URI
path segment。
mTLS Worker 身份使用 URI SAN
```text
spiffe://TRUST_DOMAIN/ENVIRONMENT/worker/WORKER_ID
```
请求中的 `worker_id` 必须与证书 URI SAN 完全一致。Session ID 不作为 TLS 身份,
也不写日志。gRPC Server 配置接收/发送消息上限、并发 Stream 上限和 keepalive
enforcement停机先停止接收新 RPC再有界 `GracefulStop`,超时后 `Stop`
Controller bootstrap 只在 `controlPlane.enabled=true` 时构造并启动 gRPC Server
并将其加入现有 lifecycle Group。控制面启用时 Redis 是必需依赖;其就绪状态纳入
Controller readiness。默认配置保持关闭避免当前部署清单把未配置证书的 `8443`
错误暴露为可运行能力。
## 7. 错误与资源边界
- 单消息最大值默认 4 MiBRuntime counter 默认最多 100,000 条。
- Proxy ID 使用与 Worker ID 相同的字符集和 128 字节上限。
- Snapshot checksum 必须是原始 SHA-256 的 32 字节值,不接受十六进制文本或其他
长度。
- ACK `error_code` 为空或匹配 `[A-Z][A-Z0-9_]{0,63}``error_message` 最长
512 字节,只做校验和丢弃。
- Counter 规范化排序后计算摘要,禁止重复 Proxy ID。
- Redis 脚本只处理有界数组,不执行全库扫描。
- 所有 gRPC 错误返回稳定 code 和低敏感度 message不返回 Redis key、Session ID、
Proxy 地址、证书内容或内部堆栈。
- 不增加无界 goroutine、无界 channel 或每请求日志。
## 8. 测试策略
测试只通过已批准的公共 seam
1. **`worker.Service` 单元测试**:注册替换、协议不兼容、负向 ACK、ACK 幂等/倒退/
checksum 错误、Runtime 心跳、Snapshot mismatch 恢复指令和存储失败。
2. **`workerruntime.ControlStore` 公共契约**:内存参考实现与 Redis Adapter 运行同一
契约,覆盖 Session TTL、ACK 原子性、报告序号和跨实例 fence。
3. **gRPC bufconn 测试**:使用生成 Client 调用三个 RPC验证 DTO、状态码、批次
限制以及两个未实现 RPC。
4. **TLS 测试**:运行时生成测试 CA/Server/Worker 证书,验证正确 SPIFFE 身份、
worker_id 不匹配、缺失客户端证书和回环明文规则。
5. **Controller bootstrap 测试**临时回环端口启动、Readiness、联动关闭和 Redis
不可用失败路径。
6. **Redis 8.2 集成测试**:通过现有 `test-redis.ps1` 验证多 Controller 共享状态、
旧 Session 隔离、TTL 过期和 ACK/Runtime 原子失败不留部分写入。
每个测试命令最长 60 秒。实现采用纵向 TDD先写一个公共行为测试并确认失败
再实现最小生产代码使其通过,然后进入下一行为。
## 9. 验收标准
1. 三个 RPC 可由生成的 gRPC Client 真实调用,并通过 Redis 共享状态。
2. 注册后未 ACK 的 Session 不能写 Runtime有效 ACK 后空报告可续期。
3. 新实例注册后,旧实例的 ACK 与 Runtime 都被拒绝。
4. Snapshot 元数据与 Session ACK 在 Redis 原子校验,不信任 Worker 自报。
5. Runtime 序号倒退、冲突、Snapshot mismatch 和存储故障具有不同恢复语义。
6. 非回环明文监听配置被拒绝mTLS Worker 身份与请求 ID 绑定。
7. `WatchSnapshots`、`ReportOutcomes` 明确返回 `Unimplemented`,文档不把本交付单元
描述为完整 Worker 控制面。
8. `go test ./...`、`go vet ./...`、`go build ./...`、protobuf 生成漂移检查和 Redis
fixture 全部通过。
9. Gateway 包及请求热路径不新增 Redis、PostgreSQL 或 gRPC Client 依赖。
## 10. 非目标
- 不实现 Snapshot payload 构建、Delta、WatchSnapshots Stream 或 Gateway 客户端。
- 不实现 Outcome 批次去重、健康 reducer 或 CheckerControlPlane。
- 不实现 Worker 维度 revoke/drain 索引;本轮 `revoke_proxy_ids` 保持为空。
- 不把 Session、ACK、Runtime 或 Snapshot 元数据写入 PostgreSQL。
- 不宣称 Gateway 进程、完整多副本控制面或 100,000 QPS 已完成。