proxy-pool/docs/superpowers/plans/2026-07-31-worker-controlplane-session-runtime.md
youfak 081e172970
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
docs: plan worker control plane session runtime
2026-07-31 10:19:30 +08:00

1282 lines
53 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# Worker Control Plane Session and Runtime Implementation Plan
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** 实现 Worker 注册、Snapshot ACK 和 Runtime 心跳三个 gRPC RPC并用 Redis 在多 Controller 副本间原子共享 Session、Snapshot 元数据和报告栅栏。
**Architecture:** protobuf 生成代码只属于传输契约,`internal/controller/worker` 提供不依赖 protobuf 的公共应用服务。`workerruntime.ControlStore` 同时由内存参考实现和 Redis Adapter 实现Controller 将独立 gRPC Runner 与现有 HTTP Runtime、Provider Supervisor 放入同一 lifecycle GroupGateway 请求热路径保持纯内存。
**Tech Stack:** Go 1.26.4、gRPC-Go v1.83.0、Protobuf-Go v1.36.11、protoc v35.0、protoc-gen-go-grpc v1.6.2、Redis 8.2、Redis Lua、mTLS/SPIFFE URI SAN、PowerShell、Docker Compose。
---
## File Structure
```text
gen/controlplane/v1/
controlplane.pb.go # protoc 生成消息类型
controlplane_grpc.pb.go # protoc 生成 gRPC Client/Server
internal/domain/workerruntime/
runtime.go # ControlStore 类型、错误和公共接口
validation.go # Service/Memory/Redis 共用规范化与摘要
validation_test.go # 公用边界规则测试
memory.go # 内存参考实现
memory_test.go # 领域行为测试
contracttest/contract.go # Memory/Redis 共用契约
internal/adapters/redisactivity/
keys.go # Worker Snapshot 元数据 key
runtime.go # ControlStore Go 适配器
scripts/runtime.lua # Session/ACK/Runtime 原子状态机
runtime_integration_test.go # Redis 8.2 契约与故障测试
internal/controller/worker/
service.go # 公共应用服务与命令/结果
service_test.go # 应用行为测试
grpc_handler.go # protobuf DTO 与 gRPC code 映射
grpc_handler_test.go # bufconn Client 测试
identity.go # SPIFFE Worker 身份校验
identity_test.go # TLS 身份测试
server.go # gRPC lifecycle Runner
server_test.go # Listener 与停机测试
internal/config/
config.go # controlPlane 配置类型
validate.go # listener、mTLS 与资源上限校验
config_test.go # 严格配置矩阵
internal/controller/bootstrap/
bootstrap.go # Worker Service/Runner 装配
infrastructure.go # Redis ControlStore 接线
bootstrap_test.go # 依赖与 lifecycle 测试
bootstrap_integration_test.go # 双存储 + gRPC fixture
scripts/
install-protoc.ps1 # 下载并校验固定 protoc 35.0
generate-proto.ps1 # 固定工具生成 Go 契约
verify-proto.ps1 # descriptor + 漂移检查
docs/configuration/reference.md # controlPlane 配置说明
docs/api/control-plane.md # 当前实现状态与恢复语义
```
### Task 1: 固定依赖并生成 gRPC 契约
**Files:**
- Modify: `go.mod`
- Modify: `go.sum`
- Create: `scripts/install-protoc.ps1`
- Create: `scripts/generate-proto.ps1`
- Modify: `scripts/verify-proto.ps1`
- Modify: `scripts/verify.ps1`
- Modify: `.github/workflows/ci.yml`
- Create: `gen/controlplane/v1/controlplane.pb.go`
- Create: `gen/controlplane/v1/controlplane_grpc.pb.go`
- [ ] **Step 1: 记录依赖与 descriptor 基线**
Run:
```powershell
go test -count=1 -timeout 60s ./...
./scripts/verify-proto.ps1
```
Expected: 全部 PASSdescriptor 非空;`gen/controlplane/v1` 尚不存在。
- [ ] **Step 2: 固定运行时与 Go tool 依赖**
Run:
```powershell
go get google.golang.org/grpc@v1.83.0 google.golang.org/protobuf@v1.36.11
go get -tool google.golang.org/protobuf/cmd/protoc-gen-go@v1.36.11
go get -tool google.golang.org/grpc/cmd/protoc-gen-go-grpc@v1.6.2
go mod tidy
```
Expected: `go.mod` 的直接依赖包含 gRPC/Protobuf并包含两个 `tool` 声明Go 版本
仍为 `1.26.0`
- [ ] **Step 3: 新增固定 protoc 安装脚本**
`scripts/install-protoc.ps1` 只支持当前 CI 和开发机的 x64 Linux/Windows
其他平台明确 throw。固定以下官方发布件与 SHA-256
```powershell
$version = "35.0"
$artifacts = @{
"linux-x86_64" = @{
File = "protoc-35.0-linux-x86_64.zip"
SHA256 = "a45cda0989c17dd950db55f6fbe1e5814c50fda08e87aa422980ac1f89dddbbc"
}
"windows-x86_64" = @{
File = "protoc-35.0-win64.zip"
SHA256 = "d1cede9e308cc3eb072392af1c02ccae4bdd3d2f374ec2970dbd8cdfdaa91363"
}
}
$baseURL = "https://github.com/protocolbuffers/protobuf/releases/download/v$version"
```
脚本将 zip 缓存到 `.tmp-proto/downloads`,先用 `Get-FileHash -Algorithm SHA256`
严格比较再 `Expand-Archive -Force``.tmp-proto/protoc-35.0`。最后执行
`protoc --version`,必须精确等于 `libprotoc 35.0`,并把可执行文件路径作为唯一
pipeline output 返回。
- [ ] **Step 4: 新增可重复生成脚本**
`scripts/generate-proto.ps1` 使用仓库内临时工具目录,不执行全局安装:
```powershell
param(
[string]$Protoc = "",
[string]$IncludePath = $env:PROTOC_INCLUDE,
[string]$OutputRoot = ""
)
$ErrorActionPreference = "Stop"
$repositoryRoot = Split-Path -Parent $PSScriptRoot
$protoRoot = Join-Path $repositoryRoot "api/proto"
$source = Join-Path $protoRoot "controlplane/v1/controlplane.proto"
$toolRoot = Join-Path $repositoryRoot ".tmp-proto/tools"
New-Item -ItemType Directory -Force -Path $toolRoot | Out-Null
if ([string]::IsNullOrWhiteSpace($Protoc)) {
$Protoc = & (Join-Path $PSScriptRoot "install-protoc.ps1")
}
$suffix = if ($IsWindows -or $env:OS -eq "Windows_NT") { ".exe" } else { "" }
$protocGenGo = Join-Path $toolRoot ("protoc-gen-go" + $suffix)
$protocGenGoGRPC = Join-Path $toolRoot ("protoc-gen-go-grpc" + $suffix)
go build -o $protocGenGo google.golang.org/protobuf/cmd/protoc-gen-go
if ($LASTEXITCODE -ne 0) { throw "build protoc-gen-go failed" }
go build -o $protocGenGoGRPC google.golang.org/grpc/cmd/protoc-gen-go-grpc
if ($LASTEXITCODE -ne 0) { throw "build protoc-gen-go-grpc failed" }
$protocCommand = Get-Command $Protoc -ErrorAction Stop
$protocVersion = & $protocCommand.Source --version
if ($protocVersion -ne "libprotoc 35.0") { throw "protoc 35.0 is required" }
if ([string]::IsNullOrWhiteSpace($IncludePath)) {
$installationRoot = Split-Path (Split-Path $protocCommand.Source -Parent) -Parent
$IncludePath = @((Join-Path $installationRoot "include"), "/usr/include", "/usr/local/include") |
Where-Object { Test-Path (Join-Path $_ "google/protobuf/timestamp.proto") } |
Select-Object -First 1
}
if ([string]::IsNullOrWhiteSpace($OutputRoot)) { $OutputRoot = $repositoryRoot }
if ([string]::IsNullOrWhiteSpace($IncludePath) -or
-not (Test-Path (Join-Path $IncludePath "google/protobuf/timestamp.proto"))) {
throw "protoc well-known type include directory was not found"
}
New-Item -ItemType Directory -Force -Path $OutputRoot | Out-Null
& $protocCommand.Source `
"--proto_path=$protoRoot" `
"--proto_path=$IncludePath" `
"--plugin=protoc-gen-go=$protocGenGo" `
"--plugin=protoc-gen-go-grpc=$protocGenGoGRPC" `
"--go_out=$OutputRoot" `
"--go_opt=module=proxy-pool" `
"--go-grpc_out=$OutputRoot" `
"--go-grpc_opt=module=proxy-pool" `
$source
if ($LASTEXITCODE -ne 0) { throw "protobuf Go generation failed" }
```
- [ ] **Step 5: 生成并编译契约**
Run:
```powershell
./scripts/generate-proto.ps1
go fmt ./gen/controlplane/v1
go mod tidy
go test -count=1 -timeout 60s ./gen/controlplane/v1
```
Expected: 两个生成文件存在package 为 `controlplanev1`,测试命令 PASS。
- [ ] **Step 6: 扩展漂移验证与 CI**
修改 `verify-proto.ps1`:保留 descriptor 校验,再生成到 `.tmp-proto/generated`
对以下两个相对路径逐字节比较:
```powershell
$generated = @(
"gen/controlplane/v1/controlplane.pb.go",
"gen/controlplane/v1/controlplane_grpc.pb.go"
)
foreach ($relative in $generated) {
$committed = Join-Path $repositoryRoot $relative
$candidate = Join-Path $outputRoot $relative
if (-not (Test-Path $committed) -or
-not [System.Linq.Enumerable]::SequenceEqual(
[System.IO.File]::ReadAllBytes($committed),
[System.IO.File]::ReadAllBytes($candidate))) {
throw "generated protobuf drift: $relative"
}
}
```
`verify-proto.ps1` 默认先调用 `install-protoc.ps1`descriptor 和 Go 生成共用返回的
固定 compiler。它把 `$outputRoot` 固定到 `.tmp-proto/generated`,通过
`& "$PSScriptRoot/generate-proto.ps1" -Protoc $Protoc -IncludePath $IncludePath -OutputRoot $outputRoot`
生成候选文件;生成失败立即 throw不使用历史候选结果。`verify.ps1` 在
存在 `protoc` 时运行扩展后的脚本。CI 新增独立 `proto` job
`ubuntu-latest` 设置 Go 后直接执行 `pwsh ./scripts/verify-proto.ps1`,不使用
apt 的浮动 `protobuf-compiler`。现有 Windows `test` job 继续由 `go test`
`go build` 验证已提交的生成包。
- [ ] **Step 7: 验证并提交生成基线**
Run:
```powershell
./scripts/verify-proto.ps1
go test -count=1 -timeout 60s ./...
git diff --check
git add go.mod go.sum scripts/install-protoc.ps1 scripts/generate-proto.ps1 scripts/verify-proto.ps1 scripts/verify.ps1 .github/workflows/ci.yml gen/controlplane/v1/controlplane.pb.go gen/controlplane/v1/controlplane_grpc.pb.go
git diff --cached --name-only
git commit -m "build: generate worker control plane grpc contract"
```
Expected: 验证 PASS提交只包含依赖、生成工具和生成代码。
### Task 2: 增加严格的 ControlPlane 配置
**Files:**
- Modify: `internal/config/config.go`
- Modify: `internal/config/validate.go`
- Modify: `internal/config/config_test.go`
- Modify: `configs/proxy-pool.yaml`
- Modify: `deploy/config/local.yaml`
- Modify: `deploy/kubernetes/base/configmap.yaml`
- Modify: `docs/configuration/reference.md`
- [ ] **Step 1: 写配置红灯测试**
`config_test.go` 新增表驱动测试,至少包含以下输入与结果:
```go
func TestValidateControlPlane(t *testing.T) {
tests := []struct {
name string
mutate func(*Config)
want string
}{
{"missing listen", func(c *Config) { c.ControlPlane = validControlPlane(); c.ControlPlane.Listen = "" }, "controlPlane listen"},
{"public plaintext", func(c *Config) { c.ControlPlane = validControlPlane(); c.ControlPlane.Listen = "0.0.0.0:8443" }, "requires mtls"},
{"short session ttl", func(c *Config) { c.ControlPlane = validControlPlane(); c.ControlPlane.SessionTTL = Duration(29 * time.Second) }, "sessionTTL"},
{"stale below heartbeat", func(c *Config) { c.ControlPlane = validControlPlane(); c.ControlPlane.MaxStaleAge = Duration(9 * time.Second) }, "maxStaleAge"},
{"invalid protocol", func(c *Config) { c.ControlPlane = validControlPlane(); c.ControlPlane.ProtocolVersion = 0 }, "protocolVersion"},
{"missing mtls files", func(c *Config) { c.ControlPlane = validMTLSControlPlane(); c.ControlPlane.TLS.CertFile = "" }, "certFile"},
{"trust domain with port", func(c *Config) { c.ControlPlane = validMTLSControlPlane(); c.ControlPlane.TLS.TrustDomain = "proxy.example:443" }, "trustDomain"},
{"environment path", func(c *Config) { c.ControlPlane = validMTLSControlPlane(); c.ControlPlane.TLS.Environment = "prod/eu" }, "environment"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
cfg := mustLoadValidConfig(t)
tt.mutate(cfg)
if err := Validate(cfg); err == nil || !strings.Contains(err.Error(), tt.want) {
t.Fatalf("Validate() error = %v, want %q", err, tt.want)
}
})
}
}
```
- [ ] **Step 2: 运行测试确认失败**
Run: `go test -count=1 -timeout 60s ./internal/config -run TestValidateControlPlane`
Expected: FAIL`Config.ControlPlane` 尚未定义。
- [ ] **Step 3: 定义配置类型**
`config.go` 增加:
```go
type Config struct {
Version int `yaml:"version"`
Defaults Defaults `yaml:"defaults"`
Security Security `yaml:"security"`
Gateway Listener `yaml:"gateway"`
Distribution Distribution `yaml:"distribution"`
Admin Listener `yaml:"admin"`
ControlPlane ControlPlane `yaml:"controlPlane"`
Metrics Metrics `yaml:"metrics"`
Storage Storage `yaml:"storage"`
Routing []Routing `yaml:"routing"`
Upstreams map[string]Upstream `yaml:"upstreams"`
}
type ControlPlane struct {
Enabled bool `yaml:"enabled"`
Listen string `yaml:"listen"`
ProtocolVersion uint32 `yaml:"protocolVersion"`
HeartbeatInterval Duration `yaml:"heartbeatInterval"`
SessionTTL Duration `yaml:"sessionTTL"`
MaxStaleAge Duration `yaml:"maxStaleAge"`
MaxMessageBytes int `yaml:"maxMessageBytes"`
MaxRuntimeCounters int `yaml:"maxRuntimeCounters"`
MaxConcurrentStreams uint32 `yaml:"maxConcurrentStreams"`
TLS ControlPlaneTLS `yaml:"tls"`
}
type ControlPlaneTLS struct {
Mode string `yaml:"mode"`
CertFile string `yaml:"certFile"`
KeyFile string `yaml:"keyFile"`
ClientCAFile string `yaml:"clientCAFile"`
TrustDomain string `yaml:"trustDomain"`
Environment string `yaml:"environment"`
}
```
- [ ] **Step 4: 实现唯一验证规则**
`Validate` 中调用 `validateControlPlane`。该函数必须执行:
```go
func validateControlPlane(item ControlPlane) error {
if !item.Enabled { return nil }
host, err := validateListenAddress("controlPlane", item.Listen)
if err != nil { return err }
if item.ProtocolVersion != 1 { return fmt.Errorf("validate controlPlane protocolVersion: must be 1") }
heartbeat, ttl, stale := item.HeartbeatInterval.Value(), item.SessionTTL.Value(), item.MaxStaleAge.Value()
if heartbeat <= 0 || ttl < 3*heartbeat { return fmt.Errorf("validate controlPlane sessionTTL: must be at least three heartbeat intervals") }
if stale < heartbeat { return fmt.Errorf("validate controlPlane maxStaleAge: must not be shorter than heartbeatInterval") }
if item.MaxMessageBytes <= 0 || item.MaxMessageBytes > 64<<20 { return fmt.Errorf("validate controlPlane maxMessageBytes: must be in [1, 67108864]") }
if item.MaxRuntimeCounters <= 0 || item.MaxRuntimeCounters > MaximumPoolSize { return fmt.Errorf("validate controlPlane maxRuntimeCounters: must be in [1, %d]", MaximumPoolSize) }
if item.MaxConcurrentStreams == 0 { return fmt.Errorf("validate controlPlane maxConcurrentStreams: must be positive") }
switch item.TLS.Mode {
case "disabled":
if isPublicHost(host) { return fmt.Errorf("validate controlPlane tls: non-loopback listen requires mtls") }
case "mtls":
if item.TLS.CertFile == "" || item.TLS.KeyFile == "" || item.TLS.ClientCAFile == "" ||
item.TLS.TrustDomain == "" || item.TLS.Environment == "" {
return fmt.Errorf("validate controlPlane tls: mtls requires certFile, keyFile, clientCAFile, trustDomain and environment")
}
if !validTrustDomain(item.TLS.TrustDomain) {
return fmt.Errorf("validate controlPlane tls.trustDomain: must be a lowercase DNS name without port")
}
if !controlPlaneIdentityPattern.MatchString(item.TLS.Environment) {
return fmt.Errorf("validate controlPlane tls.environment: must be one URI path segment")
}
default:
return fmt.Errorf("validate controlPlane tls.mode: must be disabled or mtls")
}
return nil
}
```
`validTrustDomain` 要求小写 DNS 名、总长度不超过 253每个 label 为
1..63 字节且不以 `-` 开头/结尾;禁止 userinfo、port、路径、query 和
fragment。`controlPlaneIdentityPattern` 为
`^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$`,因此 environment 不能注入额外 URI segment。
- [ ] **Step 5: 使配置测试转绿**
Run: `go test -count=1 -timeout 60s ./internal/config`
Expected: PASS包括未知字段拒绝、Redacted 深拷贝和部署配置解析。
- [ ] **Step 6: 更新真实配置与参考文档**
在默认、本地和 Kubernetes ConfigMap 的 `metrics` 前加入关闭态:
```yaml
controlPlane:
enabled: false
```
在配置参考中加入完整 loopback 示例和 mTLS 字段说明;明确默认关闭、非回环强制
mTLS、Session TTL 比例、消息/Counter 上限以及 PostgreSQL 不参与控制面状态。
- [ ] **Step 7: 验证并提交配置单元**
Run:
```powershell
go test -count=1 -timeout 60s ./internal/config ./examples/config ./deploy
go run ./deploy/tools/configcheck deploy/config/local.yaml
git diff --check
git add internal/config/config.go internal/config/validate.go internal/config/config_test.go configs/proxy-pool.yaml deploy/config/local.yaml deploy/kubernetes/base/configmap.yaml docs/configuration/reference.md
git diff --cached --name-only
git commit -m "feat: configure worker control plane listener"
```
Expected: 全部 PASS提交不包含证书或 Secret。
### Task 3: 建立 Workerruntime 公共 ControlStore 契约
**Files:**
- Modify: `internal/domain/workerruntime/runtime.go`
- Create: `internal/domain/workerruntime/validation.go`
- Create: `internal/domain/workerruntime/validation_test.go`
- Modify: `internal/domain/workerruntime/memory.go`
- Modify: `internal/domain/workerruntime/memory_test.go`
- Create: `internal/domain/workerruntime/contracttest/contract.go`
- Create: `internal/domain/workerruntime/contract_external_test.go`
- [ ] **Step 1: 写未 ACK Session 红灯测试**
新增行为测试:`OpenSession` 接受 ACK 零值Runtime 在 ACK 前返回
`ErrSnapshotMismatch`;记录引用并正向 ACK 后空 Runtime 成功。
```go
session := Session{WorkerID: "worker-a", InstanceID: "instance-a", SessionID: "session-a", ProtocolVersion: 1}
if err := store.OpenSession(ctx, session, time.Minute); err != nil { t.Fatalf("OpenSession(): %v", err) }
epoch, err := store.CurrentOwnershipEpoch(ctx)
if err != nil { t.Fatalf("CurrentOwnershipEpoch(): %v", err) }
report := Report{WorkerID: "worker-a", SessionID: "session-a", Sequence: 1, SnapshotVersion: 7, OwnershipEpoch: epoch, ObservedAt: now}
if err := store.ReplaceRuntime(ctx, report, time.Minute); !errors.Is(err, ErrSnapshotMismatch) { t.Fatalf("pre-ACK error = %v", err) }
reference := SnapshotReference{WorkerID: "worker-a", Version: 7, OwnershipEpoch: epoch, Checksum: sha256.Sum256([]byte("snapshot-7"))}
if err := store.RecordIssuedSnapshot(ctx, reference, time.Minute); err != nil { t.Fatalf("RecordIssuedSnapshot(): %v", err) }
if err := store.AcknowledgeSnapshot(ctx, SnapshotAcknowledgement{WorkerID: "worker-a", SessionID: "session-a", Reference: reference, Applied: true}, time.Minute); err != nil { t.Fatalf("AcknowledgeSnapshot(): %v", err) }
if err := store.ReplaceRuntime(ctx, report, time.Minute); err != nil { t.Fatalf("ReplaceRuntime(): %v", err) }
```
- [ ] **Step 2: 运行测试确认编译失败**
Run: `go test -count=1 -timeout 60s ./internal/domain/workerruntime`
Expected: FAIL缺少 `ControlStore` 类型和方法。
- [ ] **Step 3: 定义类型、错误和接口**
`runtime.go` 定义:
```go
var (
ErrInvalidSnapshotReference = errors.New("invalid worker snapshot reference")
ErrInvalidAcknowledgement = errors.New("invalid worker snapshot acknowledgement")
ErrSnapshotMismatch = errors.New("worker snapshot does not match acknowledged state")
ErrStaleSnapshotReference = errors.New("stale worker snapshot reference")
ErrConflictingSnapshotReference = errors.New("conflicting worker snapshot reference")
ErrStaleAcknowledgement = errors.New("stale worker snapshot acknowledgement")
)
type SnapshotReference struct {
WorkerID string
Version uint64
OwnershipEpoch uint64
Checksum [sha256.Size]byte
}
type SnapshotAcknowledgement struct {
WorkerID string
SessionID string
Reference SnapshotReference
Applied bool
ErrorCode string
}
type ControlStore interface {
CurrentOwnershipEpoch(context.Context) (uint64, error)
OpenSession(context.Context, Session, time.Duration) error
RecordIssuedSnapshot(context.Context, SnapshotReference, time.Duration) error
AcknowledgeSnapshot(context.Context, SnapshotAcknowledgement, time.Duration) error
ReplaceRuntime(context.Context, Report, time.Duration) error
}
```
`Session` 增加 `Zone`、`ProtocolVersion`、克隆后的 `Labels`
`AckedChecksum [sha256.Size]byte` 以及 `RuntimeEnabled bool``OpenSession` 要求
version、epoch 和 checksum 全部为零且 `RuntimeEnabled=false`。本 Task 先保留旧
`SessionWriter` 类型,仅为了让尚未迁移的
Redis Adapter 在中间提交仍可编译Memory 测试和新代码不再调用
`ReplaceSession`。Task 4 迁移所有 Redis 调用后立即删除该类型和 Adapter 方法,
最终树不保留兼容 seam。
`validation.go` 封装所有三层共用的边界Memory、Redis Adapter 和 Worker
Service 禁止自行复制正则或 digest 逻辑:
```go
func ValidIdentifier(string) bool
func NormalizeLabels(map[string]string) (map[string]string, error)
func NormalizeSession(Session) (Session, error)
func NormalizeSnapshotReference(SnapshotReference) (SnapshotReference, error)
func NormalizeAcknowledgement(SnapshotAcknowledgement) (SnapshotAcknowledgement, error)
func NormalizeReport(Report) (Report, [sha256.Size]byte, error)
```
`NormalizeLabels` 返回深拷贝,并执行 32 个、64/256 字节和总计 4 KiB 上限。
`NormalizeReport` 转 UTC、克隆并按 Proxy ID 排序 Counters拒绝重复 ID/负数,
再对规范 JSON 计算唯一 SHA-256 digest。Redis wire payload 可独立编码,但
Lua 比较的 digest 必须使用该公用结果。
- [ ] **Step 4: 实现内存 Session 与 epoch**
`MemoryStore` 新增 `epoch uint64`、`references map[string]memoryReference`;构造时
epoch 初始化为 `1`。`OpenSession` 总是替换同 Worker 的 Session、清除报告并复制
labels过期时间为 `now + ttl`
- [ ] **Step 5: 实现内存 Snapshot 引用与 ACK**
`RecordIssuedSnapshot` 先要求 `reference.OwnershipEpoch == CurrentOwnershipEpoch()`
再按 `(ownership_epoch, version)` 比较当前引用:新 tuple 必须严格前进;相同
tuple + 相同 checksum 为幂等续期;相同 tuple + 不同 checksum 返回
`ErrConflictingSnapshotReference`;倒退返回 `ErrStaleSnapshotReference`
ACK 必须按以下顺序执行Context → 输入 → 当前 Session/TTL → 与 Session
已 ACK 上界比较 → 当前 Reference/TTL → 与 Reference 的 tuple/checksum 比较 →
负向/正向分支。ACK tuple 落后 Session 上界或当前 Reference 均返回
`ErrStaleAcknowledgement`;超前 Reference 或同 tuple 不同 checksum 返回
`ErrSnapshotMismatch`。负向 ACK 在匹配当前 Reference 时续期、清报告并设置
`RuntimeEnabled=false`,但不推进 ACK 上界;正向 ACK 对
version/epoch/checksum 全部相同的重放幂等,并设置
`RuntimeEnabled=true`。已开启 Session 的相同正向 ACK 重放只续期 Session
不删除 Runtime 或 sequence/digest fence若相同 ACK 上界因负向 ACK 处于
`RuntimeEnabled=false`,该正向 ACK 清理任何残留 Runtime 并重新开启。
倒退上界返回 `ErrStaleAcknowledgement`。Checksum 始终用
固定数组比较,不把内容转成日志字符串。
- [ ] **Step 6: 拆分 Runtime 错误**
在内存与公共规范化逻辑中,把 Snapshot version/epoch 与 ACK 不一致改为
`ErrSnapshotMismatch`sequence 倒退继续返回 `ErrStaleReport`,相同 sequence
内容不同继续返回 `ErrConflictingReport`
- [ ] **Step 7: 写负向 ACK、替换和过期红灯测试**
新增以下公共行为:
- 负向 ACK 续期 Session、清除旧 Runtime、ACK 上界保持零。
- 已 ACK snapshot 7 后对 snapshot 8 负向 ACK延迟到达的 snapshot 7 Runtime
返回 `ErrSnapshotMismatch`;只有后续正向 ACK 才重新打开 Runtime。
- 上报 sequence 100 后重放相同正向 ACKRuntime 仍为 Fresh随后 sequence
99 仍返回 `ErrStaleReport`,证明 ACK 重放未清除 fence。
- 新 instance 注册后旧 Session 的 ACK/Runtime 返回 `ErrStaleSession`
- SnapshotReference 过期后 ACK 返回 `ErrSnapshotMismatch`
- 同一正向 ACK 的 version/epoch/checksum 幂等,旧 ACK 返回
`ErrStaleAcknowledgement`
- 同 tuple 不同 checksum 的 SnapshotReference 返回
`ErrConflictingSnapshotReference`,且当前引用不变。
- 相同 Runtime sequence 的不同内容返回 `ErrConflictingReport`
- 相同 Runtime sequence + digest 幂等重放仍续期 Session 和 Report TTL。
- [ ] **Step 8: 建立共享契约并使内存实现转绿**
`contracttest` 用封装时间推进方式的 fixture使 Memory 与 Redis 可运行
同一契约:
```go
type Fixture struct {
Store workerruntime.ControlStore
Reader workerruntime.RuntimeReader
OwnedProxy workerruntime.OwnedProxy
Advance func(time.Duration)
}
type Factory func(*testing.T) Fixture
func Run(t *testing.T, factory Factory)
```
Memory factory 的 `Store``Reader` 指向同一 `MemoryStore``OwnedProxy` 使用
`proxy-a/worker-a` 和当前 epoch`Advance` 推进受控 `time.Time`Redis integration
factory 先通过公共 Activity/Ownership API 创建并分配同名代理,把 assignment
写入 `OwnedProxy`,它的
`Advance` 使用有界 `time.Sleep`。`Run` 对任意 `ControlStore` 运行上述行为,
每个 subtest 都调用 factory 获得独立 Store/namespace避免状态串扰。
外部测试用 `NewMemoryStore` 注册该契约。
Run: `go test -count=1 -timeout 60s ./internal/domain/workerruntime`
Expected: PASS。
- [ ] **Step 9: 提交领域契约**
Run:
```powershell
gofmt -w internal/domain/workerruntime
go test -count=1 -timeout 60s ./...
git diff --check
git add internal/domain/workerruntime/runtime.go internal/domain/workerruntime/validation.go internal/domain/workerruntime/validation_test.go internal/domain/workerruntime/memory.go internal/domain/workerruntime/memory_test.go internal/domain/workerruntime/contracttest/contract.go internal/domain/workerruntime/contract_external_test.go
git diff --cached --name-only
git commit -m "feat: define worker control store contract"
```
Expected staged files 精确等于上述 7 个路径,且全仓测试 PASS。
### Task 4: 用 Redis Lua 原子实现 ControlStore
**Files:**
- Modify: `internal/adapters/redisactivity/keys.go`
- Modify: `internal/adapters/redisactivity/adapter_test.go`
- Modify: `internal/adapters/redisactivity/scripts.go`
- Modify: `internal/adapters/redisactivity/runtime.go`
- Modify: `internal/adapters/redisactivity/scripts/runtime.lua`
- Modify: `internal/adapters/redisactivity/scripts/ownership.lua`
- Modify: `internal/adapters/redisactivity/ownership_integration_test.go`
- Modify: `internal/adapters/redisactivity/runtime_integration_test.go`
- Modify: `internal/adapters/redisactivity/capacity_integration_test.go`
- [ ] **Step 1: 让 Redis Adapter 运行共享契约并确认失败**
`runtime_integration_test.go` 的 Redis 8.2 fixture 中断言
`redisactivity.Adapter` 实现 `workerruntime.ControlStore`。每次 factory 调用生成唯一
namespace仅在验证跨 Controller fence 的专项 subtest 内,两个 Adapter
共享该 subtest 的 namespace。工厂传入 `workerruntime/contracttest.Run`
Run:
```powershell
./scripts/test-redis.ps1
```
Expected: FAILAdapter 缺少新方法。
- [ ] **Step 2: 增加同槽位 Snapshot key**
`keyspace` 新增:
```go
workerSnapshots string
workerSnapshotExpiry string
```
值分别为 `prefix + ":worker-snapshots"`
`prefix + ":worker-snapshot-expiry"`,继续共享 `pp:{activity}:` hash tag。
`adapter_test.go` 更新完整 keyspace 断言,并对 Session、Snapshot、Runtime、
Owner 和 epoch 全部 key 执行 Redis Cluster slot 一致性检查。
- [ ] **Step 3: 增加明确脚本状态与 Wire DTO**
新增 `scriptSnapshotMismatch = "snapshot_mismatch"`
`scriptStaleAcknowledgement = "stale_acknowledgement"`。Wire DTO 中所有 uint64
继续使用十进制字符串checksum 使用 64 字符小写 hexGo 边界仍是 `[32]byte`
`runtimeSessionWire` 新增 `ackedChecksum``runtimeEnabled`:未 ACK Session 表示为
version=`"0"`、epoch=`"0"`、checksum=`""`、runtimeEnabled=`false`;已 ACK 上界
必须为两个正数和 64 位 checksum。`runtimeEnabled=false` 也可与非零旧 ACK
上界共存,表示新 Snapshot 负向 ACK 后的 fail-closed 状态。
- [ ] **Step 4: 扩展 Lua 当前 epoch 与 Session 注册**
`current_epoch` 分支执行 `SETNX epoch_key '1'`、`PERSIST epoch_key` 后返回
当前十进制 epoch。`ownership.lua` 的 Assign/Renew 不再对全局 `epoch_key`
执行 `PEXPIREAT`Assign 在 `INCR` 后执行 `PERSIST`Renew 也对现有 key
执行 `PERSIST`,以迁移历史带 TTL 的 namespace。该 namespace 只保留这一个
持久单调标量,避免 TTL 到期后 epoch 回退造成 ABA。
`open_session` 分支必须:
1. 校验 ACK version/epoch 都是字符串 `"0"`、`ackedChecksum == ""` 且
`runtimeEnabled == false`
2. 删除该 Worker runtime 和 runtime expiry。
3. 写入新 Session 与服务端时间 expiry。
4. 不读取客户端时间决定 TTL。
- [ ] **Step 5: 实现 SnapshotReference 原子写入**
`record_snapshot` 校验 version/epoch 正数、checksum 为 64 位小写 hex
epoch 精确等于当前持久 `epoch_key`。它按 `(epoch, version)` 实现与 Memory
相同的严格单调规则:同 tuple 同 checksum 幂等续期,同 tuple 不同 checksum
返回 `conflict`,倒退返回 `stale`,两种失败都不修改 Hash/ZSET。
成功时写入 Worker Hash 并以 Redis `TIME` 设置 ZSET expiry。每次 Lua 调用先用
`ZRANGEBYSCORE ... LIMIT 0 cleanup_limit` 有界删除过期 Snapshot Hash/ZSET 条目;
新引用覆盖旧引用。
- [ ] **Step 6: 实现 ACK 原子状态机**
`ack_snapshot` 在一个 Lua 调用中完成:
```lua
-- 校验当前 session 与 snapshot reference 后:
if applied == false then
redis.call('HDEL', runtime_key, worker_id)
redis.call('ZREM', runtime_expiry_key, worker_id)
session.runtimeEnabled = false
session.expiresAtMs = now + ttl_ms
redis.call('HSET', sessions_key, worker_id, cjson.encode(session))
redis.call('ZADD', session_expiry_key, session.expiresAtMs, worker_id)
return cjson.encode({status = 'ok'})
end
-- 先判断:倒退返回 stale_acknowledgement。
-- 相同正向 ACK + runtimeEnabled=true只续期 Session不修改 Runtime/fence。
-- 相同正向 ACK + runtimeEnabled=false清理 Runtime并重新开启。
-- 严格前进正向 ACK清理 Runtime、保存 version/epoch/checksum 并开启。
```
引用不匹配、引用过期或未发布统一返回 `snapshot_mismatch`Session 缺失/替换返回
`unavailable`Go 映射为 `ErrStaleSession`
- [ ] **Step 7: 修改 Runtime 分支错误分类**
Session 未 ACK、`runtimeEnabled=false` 或报告 version/epoch 不等于 ACK 时返回
`snapshot_mismatch`。sequence
倒退仍返回 `stale`digest 冲突仍返回 `conflict`。相同 sequence + digest
的幂等重放不能直接 return它要用当前 Redis `TIME` 重写 Session/Report
`expiresAtMs` 与两个 expiry ZSET score。新 sequence 成功时完整替换报告,
并同时续期 Session 与报告 TTL。
- [ ] **Step 8: 实现 Go Adapter 方法与映射**
`runtime.go` 实现五个 `ControlStore` 方法;所有调用经过 Context、nil Adapter、
固定长度、数量上限和 TTL 校验。`CurrentOwnershipEpoch` 解析字符串时拒绝零值、负数
和超出 uint64 的响应。
字段规范化和 Runtime digest 一律调用 `workerruntime` 公用方法Adapter 只追加
`MaxRuntimeCounters` 配置上限和 wire 编码。
`runRuntime` 对 Lua 始终传入以下同槽 key 顺序,所有分支共用该索引:
```go
[]string{
a.keys.workerSessions,
a.keys.workerSessionExpiry,
a.keys.workerSnapshots,
a.keys.workerSnapshotExpiry,
a.keys.workerRuntime,
a.keys.workerRuntimeExpiry,
a.keys.owners,
a.keys.epoch,
}
```
`capacity_integration_test.go` 不再调用已删除的 `ReplaceSession`;共用测试 helper
`OpenSession -> RecordIssuedSnapshot -> AcknowledgeSnapshot` 建立已 ACK Session
再执行原有 capacity 断言。
迁移 `runtime_integration_test.go``capacity_integration_test.go` 后,删除
`workerruntime.SessionWriter` 以及 Redis Adapter 的 `ReplaceSession`;保留名称不变的
`ReportWriter.ReplaceRuntime``RuntimeReader`,因为它们仍是有效的窄端口。
- [ ] **Step 9: 增加真实 Redis 并发与故障测试**
新增:
- 两个 Adapter 同 namespace 注册同 Worker新 Session 隔离旧 Session。
- 100 个并发相同 ACK 全部幂等成功,最终只有一个 ACK 上界。
- 相同 Snapshot tuple 换 checksum 返回冲突,旧 Reference/Session/Runtime 完全不变。
- Snapshot 8 负向 ACK 后,延迟的 snapshot 7 Runtime 在两个 Adapter 上都被
`snapshot_mismatch` 拒绝,不得重建已清除的 sequence/report。
- sequence 100 后幂等重放当前正向 ACK旧 Runtime 仍 Freshsequence 99
仍被拒绝,证明 ACK 重放未清除 Runtime fence。
- 正向 ACK 校验失败后 Session/Runtime/Reference 均无部分写入。
- 100ms Session/Snapshot TTL 到期后 fail-closed。
- 人为给 epoch key 设置短 TTL 后调用 `CurrentOwnershipEpoch`,断言 key 变为
persistent后续分配的 epoch 严格增长,旧 Runtime 对新归属仍为 not Fresh。
- Redis 断开时五个方法返回 Store unavailable 类错误而非 panic。
- [ ] **Step 10: 验证并提交 Redis 单元**
Run:
```powershell
go test -count=1 -timeout 60s ./internal/adapters/redisactivity ./internal/domain/workerruntime
./scripts/test-redis.ps1
go test -count=1 -timeout 60s ./...
git diff --check
git add internal/domain/workerruntime/runtime.go internal/adapters/redisactivity/keys.go internal/adapters/redisactivity/adapter_test.go internal/adapters/redisactivity/scripts.go internal/adapters/redisactivity/runtime.go internal/adapters/redisactivity/scripts/runtime.lua internal/adapters/redisactivity/scripts/ownership.lua internal/adapters/redisactivity/ownership_integration_test.go internal/adapters/redisactivity/runtime_integration_test.go internal/adapters/redisactivity/capacity_integration_test.go
git diff --cached --name-only
git commit -m "feat: persist worker sessions and snapshot acknowledgements"
```
Expected: 内存/Redis 共用契约与真实 Redis 8.2 集成测试 PASS。
### Task 5: 实现 Worker 公共应用服务
**Files:**
- Create: `internal/controller/worker/service.go`
- Create: `internal/controller/worker/service_test.go`
- [ ] **Step 1: 写注册服务红灯测试**
使用记录型 `ControlStore` 与确定性 Session ID 生成器,验证:协议版本 1、响应
返回 heartbeat/max stale age、`OpenSession` 收到配置的 SessionTTL、Session ACK
version/epoch/checksum 零值以及 Current epoch 透传。Proto 响应不断言不存在的
SessionTTL 字段。
- [ ] **Step 2: 运行测试确认失败**
Run: `go test -count=1 -timeout 60s ./internal/controller/worker`
Expected: FAILpackage 或 `NewService` 尚不存在。
- [ ] **Step 3: 定义公共命令、结果和选项**
`service.go` 定义:
```go
type RegisterCommand struct {
WorkerID string
InstanceID string
Zone string
ProtocolVersion uint32
Labels map[string]string
}
type SnapshotAcknowledgement struct {
WorkerID string
SessionID string
Version uint64
OwnershipEpoch uint64
Checksum []byte
Applied bool
ErrorCode string
ErrorMessage string
}
type Registration struct {
WorkerID string
SessionID string
OwnershipEpoch uint64
HeartbeatInterval time.Duration
MaxStaleAge time.Duration
}
type RuntimeDecision struct {
AcceptedOwnershipEpoch uint64
RequireFullSnapshot bool
}
type Options struct {
ProtocolVersion uint32
HeartbeatInterval time.Duration
SessionTTL time.Duration
MaxStaleAge time.Duration
MaxRuntimeCounters int
SessionID func() (string, error)
}
type Service interface {
Register(context.Context, RegisterCommand) (Registration, error)
Acknowledge(context.Context, SnapshotAcknowledgement) error
ReportRuntime(context.Context, workerruntime.Report) (RuntimeDecision, error)
}
var (
ErrInvalidCommand = errors.New("invalid worker control command")
ErrProtocolVersion = errors.New("unsupported worker protocol version")
ErrUnavailable = errors.New("worker control service unavailable")
)
```
默认 Session ID 生成器使用 `crypto/rand.Read` 读取 16 字节并编码成 32 字符小写
hex失败必须返回内部错误不降级到时间戳或伪随机数。
`NewService` 返回实现上述接口的非导出类型,调用方不依赖具体结构。
服务的 ID、labels、Session、ACK 和 Runtime 校验调用 `workerruntime` 公用规范化
方法,并把领域的 invalid sentinel 统一包装为 `ErrInvalidCommand`
- [ ] **Step 4: 实现 Register**
验证 ID pattern、Zone、32 个标签和 4 KiB 总长度;读取 epoch、生成 Session、调用
`OpenSession`。存储错误对上层只暴露稳定的 `ErrUnavailable`,原始错误只用于
`errors.Is` 链,不包含 Session ID。
- [ ] **Step 5: 写并实现 ACK 行为**
测试并实现32 字节 checksum、error code pattern、512 字节 message 上限、正向/
负向 ACK 调用 ControlStoreService 校验 `ErrorMessage` 的字节长度后,构造
`workerruntime.SnapshotAcknowledgement`时只传入 `ErrorCode`,不向 Store 传递
自由文本。
- [ ] **Step 6: 写并实现 Runtime 决策**
`ReportRuntime``ErrSnapshotMismatch` 返回:
```go
RuntimeDecision{AcceptedOwnershipEpoch: currentEpoch, RequireFullSnapshot: true}, nil
```
成功返回当前 epoch 和 `RequireFullSnapshot:false`stale/conflict/session/storage 错误
保持不同 sentinel供 gRPC 层映射。
- [ ] **Step 7: 覆盖边界与并发**
测试协议不兼容、Session ID 生成失败、重复 Proxy ID、Counter 超限、Context 取消、
空 Runtime 心跳、Store 不可用和 100 个并发注册最终只有最后 Session 有效。
- [ ] **Step 8: 验证并提交服务**
Run:
```powershell
gofmt -w internal/controller/worker
go test -count=1 -timeout 60s ./internal/controller/worker
git diff --check
git add internal/controller/worker/service.go internal/controller/worker/service_test.go
git diff --cached --name-only
git commit -m "feat: handle worker sessions and runtime reports"
```
### Task 6: 实现 gRPC Handler、SPIFFE 身份与 Server Runner
**Files:**
- Create: `internal/controller/worker/grpc_handler.go`
- Create: `internal/controller/worker/grpc_handler_test.go`
- Create: `internal/controller/worker/identity.go`
- Create: `internal/controller/worker/identity_test.go`
- Create: `internal/controller/worker/server.go`
- Create: `internal/controller/worker/server_test.go`
- [ ] **Step 1: 写 bufconn 注册红灯测试**
使用生成的 `NewWorkerControlPlaneClient` 连接 bufconn调用 `RegisterWorker` 并断言
DTO、Duration 和 Session测试通过真实 gRPC 编解码,不直接调用 Handler 方法。
- [ ] **Step 2: 运行测试确认失败**
Run: `go test -count=1 -timeout 60s ./internal/controller/worker -run TestGRPC`
Expected: FAIL缺少 gRPC Handler。
- [ ] **Step 3: 实现 Handler 和 DTO 映射**
`GRPCHandler` 嵌入生成的 `UnimplementedWorkerControlPlaneServer`,只覆盖三个 RPC。
`AcknowledgeSnapshot` 使用 `emptypb.Empty`timestamp 用 `AsTime` 前先调用
`CheckValid`Runtime uint32 安全转换为 int64。
Handler 只依赖公共接口,并在进入 Service 前完成身份校验:
```go
type IdentityAuthorizer interface {
Authorize(context.Context, string) error
}
type GRPCHandler struct {
controlplanev1.UnimplementedWorkerControlPlaneServer
service Service
identity IdentityAuthorizer
}
```
三个 RPC 都以请求中的 `worker_id` 调用 `Authorize`;身份失败固定映射
`codes.PermissionDenied`,不把证书主体或 URI 写入 gRPC message。
- [ ] **Step 4: 固定 gRPC 状态码映射**
实现唯一映射:
```go
switch {
case errors.Is(err, context.Canceled):
return status.Error(codes.Canceled, "worker control request canceled")
case errors.Is(err, context.DeadlineExceeded):
return status.Error(codes.DeadlineExceeded, "worker control request deadline exceeded")
case errors.Is(err, ErrInvalidCommand),
errors.Is(err, workerruntime.ErrInvalidReport),
errors.Is(err, workerruntime.ErrInvalidAcknowledgement),
errors.Is(err, workerruntime.ErrInvalidSnapshotReference):
return status.Error(codes.InvalidArgument, "invalid worker control request")
case errors.Is(err, ErrProtocolVersion):
return status.Error(codes.FailedPrecondition, "unsupported worker protocol version")
case errors.Is(err, workerruntime.ErrStaleSession):
return status.Error(codes.FailedPrecondition, "worker session is stale")
case errors.Is(err, workerruntime.ErrSnapshotMismatch):
return status.Error(codes.FailedPrecondition, "worker snapshot does not match issued snapshot")
case errors.Is(err, workerruntime.ErrStaleAcknowledgement):
return status.Error(codes.Aborted, "worker snapshot acknowledgement is stale")
case errors.Is(err, workerruntime.ErrStaleReport):
return status.Error(codes.Aborted, "worker runtime sequence is stale")
case errors.Is(err, workerruntime.ErrConflictingReport):
return status.Error(codes.AlreadyExists, "worker runtime sequence conflicts")
default:
return status.Error(codes.Unavailable, "worker control plane unavailable")
}
```
- [ ] **Step 5: 验证两个未实现 RPC**
通过生成 Client 调用 `WatchSnapshots`,断言首次 `Recv()` 返回
`codes.Unimplemented`;调用 `ReportOutcomes`后执行 `CloseAndRecv()`,断言
`codes.Unimplemented`。不要添加空成功响应。
- [ ] **Step 6: 写 SPIFFE 身份红灯测试**
测试 URI SAN `spiffe://proxy.example/prod/worker/worker-a` 可授权 `worker-a`,以下输入
返回 PermissionDenied缺失 TLSInfo、缺失 URI SAN、错误 trust domain/environment、
Checker 类型、请求 worker ID 不一致、多个冲突 Worker URI。
- [ ] **Step 7: 实现 IdentityAuthorizer**
`peer.FromContext` 取得 `credentials.TLSInfo`,只读取已验证链叶子证书 URI SAN。
使用 URL path segment 精确比较,不做 substring 或前缀授权。plaintext mode 注入
`AllowLoopbackIdentity`,只在配置已验证为回环监听时构造。
- [ ] **Step 8: 写 Server Runner 红灯测试**
测试 `Run(ctx)``127.0.0.1:0` 启动、Context 取消后有界退出、占用端口时报错;
mTLS 测试运行时生成 CA/Server/Worker 证书,不提交 `.pem``.key`
另用预绑定 `net.Listener` 调用 `Serve(ctx, listener)`,通过
`listener.Addr()` 发起真实 Client 请求。
真实 gRPC Client 测试还覆盖:
- 超过 `MaxMessageBytes` 的 Register 请求返回 `codes.ResourceExhausted`
- 正确 CA + Worker URI SAN 完成 mTLS 握手并调用成功。
- 缺失客户端证书、错误 CA 在握手期失败。
- 合法证书但 Worker ID 不一致时 RPC 返回 `codes.PermissionDenied`
- 服务依赖返回 `context.Canceled` / `context.DeadlineExceeded`Client 分别收到
`codes.Canceled` / `codes.DeadlineExceeded`
- [ ] **Step 9: 实现 gRPC Server**
`NewServer` 配置:
- `grpc.MaxRecvMsgSize` / `MaxSendMsgSize`
- `grpc.MaxConcurrentStreams`
- keepalive enforcement 最小 ping 10 秒、无活动连接禁止 ping
- mTLS 使用 `tls.RequireAndVerifyClientCert` 和独立 Client CA pool
- Context 取消后 goroutine 执行 `GracefulStop`,到 ShutdownTimeout 后调用 `Stop`
```go
type ServerOptions struct {
ShutdownTimeout time.Duration
}
func DefaultServerOptions() ServerOptions {
return ServerOptions{ShutdownTimeout: 15 * time.Second}
}
```
`NewServer` 对零值应用 15 秒默认值,对负值返回 `ErrInvalidServer`
Server 同时提供生产绑定与可测试绑定:
```go
func (server *Server) Run(ctx context.Context) error
func (server *Server) Serve(ctx context.Context, listener net.Listener) error
```
`Run` 使用已验证的 `controlPlane.listen` 创建 TCP Listener然后委托 `Serve`
`Serve` 拥有传入 Listener 的关闭责任,且拒绝 nil Context/Listener。
- [ ] **Step 10: 验证并提交传输单元**
Run:
```powershell
gofmt -w internal/controller/worker
go test -count=1 -timeout 60s ./internal/controller/worker
git diff --check
git add internal/controller/worker/grpc_handler.go internal/controller/worker/grpc_handler_test.go internal/controller/worker/identity.go internal/controller/worker/identity_test.go internal/controller/worker/server.go internal/controller/worker/server_test.go
git diff --cached --name-only
git commit -m "feat: serve worker control plane grpc"
```
### Task 7: 接入 Controller Infrastructure 与 Lifecycle
**Files:**
- Modify: `internal/controller/bootstrap/infrastructure.go`
- Modify: `internal/controller/bootstrap/infrastructure_test.go`
- Modify: `internal/controller/bootstrap/bootstrap.go`
- Modify: `internal/controller/bootstrap/bootstrap_test.go`
- Modify: `internal/controller/bootstrap/bootstrap_integration_test.go`
- [ ] **Step 1: 写 Infrastructure 红灯测试**
验证 `controlPlane.enabled=true` 时,即使 Distribution/Admin/Provider 都关闭,也会
打开 Redis、构造 `workerruntime.ControlStore` 和 Redis readinessRedis URL 缺失或
Ping 失败时启动失败。
- [ ] **Step 2: 修改 ports 与 Redis 打开条件**
`ports` 新增:
```go
workerControl workerruntime.ControlStore
```
Redis 条件改为:
```go
configuration.Distribution.Enabled || configuration.Admin.Enabled ||
configuration.ControlPlane.Enabled || providersEnabled
```
同一个 `redisactivity.Adapter` 同时赋给 `opened.activity`
`opened.workerControl`,不创建第二客户端或第二 namespace。
Adapter 的 Runtime 批次上限由公用 helper 取当前容量与 ControlPlane 上限
的较大值,避免 control-plane-only 配置退化到 `1`
```go
func runtimeCounterCapacity(configuration *config.Config) int {
capacity := credentialCapacity(configuration)
if configuration.ControlPlane.Enabled && configuration.ControlPlane.MaxRuntimeCounters > capacity {
return configuration.ControlPlane.MaxRuntimeCounters
}
return capacity
}
```
`redisactivity.Options.MaxRuntimeCounters` 只使用该 helper。表驱动测试覆盖
control-plane-only `100000`、Provider/Gateway 容量较大以及 ControlPlane 关闭三种情况。
- [ ] **Step 3: 写 Bootstrap 红灯测试**
使用记录型 worker runner factory验证开启 controlPlane 时构造 Service 与 Runner
缺少 ControlStore、Runner factory 返回 nil、构造错误均返回 `ErrStartup`
- [ ] **Step 4: 接入 Worker Service 与 Runner**
新增独立 `workerRuntimeFactory` seam生产实现调用 `worker.NewServer`。将
Bootstrap 内部签名固定为:
```go
type workerRuntimeFactory interface {
New(config.ControlPlane, worker.Service) (controllerRunner, error)
}
func run(
ctx context.Context,
options Options,
infrastructure infrastructure,
factory runtimeFactory,
workerFactory workerRuntimeFactory,
) error
```
`Run` 传入 `productionRuntimeFactory{}``productionWorkerRuntimeFactory{}`;所有现有
Bootstrap 测试显式传入记录型 Worker factory。Bootstrap 使用配置构造
`worker.Options`,并把 Worker Runner 追加到 `runners`
```go
if loaded.Value.ControlPlane.Enabled {
if nilInterface(opened.workerControl) { return errors.Join(ErrStartup, ErrInvalidOptions) }
controlPlane := loaded.Value.ControlPlane
service, err := worker.NewService(opened.workerControl, worker.Options{
ProtocolVersion: controlPlane.ProtocolVersion,
HeartbeatInterval: controlPlane.HeartbeatInterval.Value(),
SessionTTL: controlPlane.SessionTTL.Value(),
MaxStaleAge: controlPlane.MaxStaleAge.Value(),
MaxRuntimeCounters: controlPlane.MaxRuntimeCounters,
})
if err != nil { return fmt.Errorf("%w: build Worker service: %w", ErrStartup, err) }
runner, err := workerFactory.New(controlPlane, service)
if err != nil { return fmt.Errorf("%w: build Worker gRPC runtime: %w", ErrStartup, err) }
if nilInterface(runner) { return errors.Join(ErrStartup, ErrInvalidOptions) }
runners = append(runners, runner)
}
```
`productionWorkerRuntimeFactory.New` 调用
`worker.NewServer(controlPlane, service, worker.DefaultServerOptions())`,不在 Bootstrap 复制
TLS/listener 构造逻辑。
- [ ] **Step 5: 把控制面纳入 Readiness**
`selectMetricsReadiness` 在 controlPlane 启用时至少检查 Redis。已有 Distribution 或
Provider 时继续复用 Redis readinessAdmin-only 规则保持不变。测试 Redis 失败使
`/readyz` 返回 503`/livez` 仍返回 200。
- [ ] **Step 6: 测试生命周期联动**
启动 HTTP、gRPC、Provider 三个记录型 Runner任一 Runner 返回错误时,其余 Context
都被取消;正常取消时三者全部有界退出,资源关闭只执行一次。
- [ ] **Step 7: 扩展真实 Controller fixture**
`bootstrap_integration_test.go` 用 loopback plaintext controlPlane 和真实 Redis
启动 Controller。集成测试的 `workerRuntimeFactory` 预绑定
`127.0.0.1:0`,返回调用 `server.Serve(ctx, listener)` 的 Runner并向测试暴露
`listener.Addr()`。通过该真实地址创建生成 gRPC Client 完成 Register。测试另用
`redisactivity.New` 连接同一 fixture Redis 和同一 namespace先通过现有
`UpsertFetched` 和两次 `ApplyHealth` 公共方法把 `proxy-a` 推进到
`checking -> available`,再调用 `Assign(ctx, now, "proxy-a", "worker-a", ttl)` 获得
权威 ownership epoch并用该 epoch 调用 `RecordIssuedSnapshot`。通过 gRPC 完成
ACK 和 counters 为空的 Runtime 后,用测试 Adapter 调用:
```go
ReadRuntime(ctx, []workerruntime.OwnedProxy{{
ProxyID: "proxy-a", WorkerID: "worker-a", OwnershipEpoch: assignment.Epoch,
}})
```
断言返回 `Fresh:true`、`Active:0`、`Reserved:0`,以公共 seam 证明空报告
既是完整稀疏替换,也完成 Session 心跳续期。
- [ ] **Step 8: 验证并提交 Bootstrap 单元**
Run:
```powershell
gofmt -w internal/controller/bootstrap
go test -count=1 -timeout 60s ./internal/controller/bootstrap
./scripts/test-controller.ps1
git diff --check
git add internal/controller/bootstrap/infrastructure.go internal/controller/bootstrap/infrastructure_test.go internal/controller/bootstrap/bootstrap.go internal/controller/bootstrap/bootstrap_test.go internal/controller/bootstrap/bootstrap_integration_test.go
git diff --cached --name-only
git commit -m "feat: run worker grpc control plane in controller"
```
### Task 8: 更新契约文档、实施进度与质量门禁
**Files:**
- Modify: `docs/api/control-plane.md`
- Modify: `docs/development/implementation-plan.md`
- Modify: `docs/requirements/traceability.md`
- Modify: `docs/requirements/completion-audit.md`
- Modify: `docs/adr/005-redis-activity-pool.md`
- Modify: `README.md`
- [ ] **Step 1: 更新准确状态**
文档必须明确:
- Register/ACK/Runtime gRPC 已实现并有 Redis 8.2 证据。
- Runtime 空报告承担心跳Session/ACK/report 使用 Redis 服务端 TTL。
- 负向 ACK 关闭 Session Runtime 写入栅栏,直到新正向 ACK延迟旧报告不可恢复。
- WatchSnapshots/ReportOutcomes 仍返回 Unimplemented。
- Gateway 命令、Snapshot payload/stream、Outcome 和 Checker 仍未闭环。
- 100,000 QPS 仍是未验证设计目标。
- [ ] **Step 2: 更新实施计划检查项**
只勾选 Task 10 的“Worker heartbeat receiving path and session lifecycle”。不得勾选
完整 Snapshot、Gateway 客户端、Checker 或部署拓扑任务。完成数从
`52/74 (70.3%)` 更新为 `53/74 (71.6%)`README、完成度审计和实施计划
使用同一统计日期 `2026-07-31`
- [ ] **Step 3: 校准 Kubernetes 清单**
Controller `8443` 端口和 NetworkPolicy 保留,但 ConfigMap 继续
`controlPlane.enabled:false`,直到环境 Overlay 提供 mTLS 文件、trust domain 与
environment。本步只读核对 `controller.yaml``networkpolicy.yaml`,不修改、
不暂存;文档说明静态端口不代表控制面默认启用。
- [ ] **Step 4: 执行定向测试**
Run:
```powershell
go test -count=1 -timeout 60s ./internal/domain/workerruntime ./internal/adapters/redisactivity ./internal/controller/worker ./internal/controller/bootstrap ./internal/config ./docs ./deploy
./scripts/test-redis.ps1
./scripts/test-controller.ps1
./scripts/verify-proto.ps1
```
Expected: 全部 PASS。
- [ ] **Step 5: 执行全仓质量门禁**
Run:
```powershell
./scripts/verify.ps1
go test -count=1 -timeout 60s ./...
go vet ./...
go build ./...
git diff --check
```
Expected: 全部 PASS本机 CGO 关闭时明确记录 race 由 Linux CI 执行。
- [ ] **Step 6: 执行边界扫描**
Run:
```powershell
$gatewayMatches = rg -n "redis|postgres|grpc" internal/gateway
if ($LASTEXITCODE -gt 1) { throw "gateway boundary scan failed" }
$gatewayMatches
rg -n "session_id|worker_id|proxy_id|proxy_ip" internal/platform/metrics internal/controller/worker
git status --short
```
Expected: Gateway 没有新增存储或 gRPC Client 导入;指标没有高基数标签;工作树只含
本任务预期文件。
- [ ] **Step 7: 提交并推送**
Run:
```powershell
git add README.md docs/api/control-plane.md docs/development/implementation-plan.md docs/requirements/traceability.md docs/requirements/completion-audit.md docs/adr/005-redis-activity-pool.md
git diff --cached --name-only
git commit -m "docs: record worker control plane delivery"
git push origin build/proxy-pool-architecture
git rev-list --left-right --count origin/build/proxy-pool-architecture...HEAD
```
Expected: 最终输出 `0 0`,远端分支与本地同步。