From 081e1729706eb18162a94dac24dac04b7d13b9e7 Mon Sep 17 00:00:00 2001 From: youfak Date: Fri, 31 Jul 2026 10:19:30 +0800 Subject: [PATCH] docs: plan worker control plane session runtime --- ...-31-worker-controlplane-session-runtime.md | 1281 +++++++++++++++++ ...ker-controlplane-session-runtime-design.md | 46 +- 2 files changed, 1317 insertions(+), 10 deletions(-) create mode 100644 docs/superpowers/plans/2026-07-31-worker-controlplane-session-runtime.md diff --git a/docs/superpowers/plans/2026-07-31-worker-controlplane-session-runtime.md b/docs/superpowers/plans/2026-07-31-worker-controlplane-session-runtime.md new file mode 100644 index 0000000..6234fb0 --- /dev/null +++ b/docs/superpowers/plans/2026-07-31-worker-controlplane-session-runtime.md @@ -0,0 +1,1281 @@ +# 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 Group,Gateway 请求热路径保持纯内存。 + +**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: 全部 PASS,descriptor 非空;`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 后重放相同正向 ACK,Runtime 仍为 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: FAIL,Adapter 缺少新方法。 + +- [ ] **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 字符小写 hex,Go 边界仍是 `[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 仍 Fresh,sequence 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: FAIL,package 或 `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 调用 ControlStore;Service 校验 `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 readiness;Redis 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 readiness;Admin-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`,远端分支与本地同步。 diff --git a/docs/superpowers/specs/2026-07-31-worker-controlplane-session-runtime-design.md b/docs/superpowers/specs/2026-07-31-worker-controlplane-session-runtime-design.md index 99f6ea6..0d6d97a 100644 --- a/docs/superpowers/specs/2026-07-31-worker-controlplane-session-runtime-design.md +++ b/docs/superpowers/specs/2026-07-31-worker-controlplane-session-runtime-design.md @@ -95,9 +95,15 @@ Cluster slot。 - `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 有以下状态: @@ -116,6 +122,13 @@ Session 有以下状态: | 空或非空 Runtime +--------------------> [续期] | + | 对新 Snapshot 的负向 ACK + v +[保留旧 ACK,Runtime 已禁用] + | + | 后续有效正向 ACK + +--------------------> [已 ACK,可接收 Runtime] + | | 新 instance 注册 / TTL 过期 v [失效] @@ -125,7 +138,8 @@ Session 有以下状态: 1. 每次成功注册都生成新的 128-bit 加密随机 Session ID。 2. 同一 `worker_id` 再注册会替换旧 Session,并删除旧 Runtime 报告。 -3. 注册后的 `acked_snapshot_version` 和 `acked_ownership_epoch` 均为零。 +3. 注册后的 `acked_snapshot_version`、`acked_ownership_epoch` 和 + `acked_checksum` 均为零值,`runtime_enabled=false`。 4. Runtime 只有在 Session 已完成正向 ACK 后才可接受。 5. 空 Runtime 报告表示完整稀疏计数全部为零,同时承担心跳并续期 Session。 6. Session 到期后,旧实例的 ACK 和 Runtime 均失败;使用该报告计算容量时保持 @@ -134,7 +148,9 @@ Session 有以下状态: `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;禁止重复键、空键和首尾空白。 +4 KiB;禁止空键和首尾空白。`labels` 在 protobuf 契约中是 map,重复 wire key +在进入 Service 前已按 protobuf map 语义归一化为最后一值;Service 只验证 +解码后的唯一键集。 ## 5. RPC 行为 @@ -155,12 +171,19 @@ Session 有以下状态: ACK 操作在 Redis Lua 的一个原子边界内完成: 1. Session 必须存在、未过期,并匹配 worker/session。 -2. version、ownership epoch 和 checksum 必须匹配当前已发布 SnapshotReference。 -3. `applied=true` 时,ACK 只能严格前进或完全相同地幂等重放;成功后续期 Session - 并清除旧 Runtime。 +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,也不持久化自由文本错误;下一轮 Snapshot Stream 据此继续 - 发送完整 Snapshot。 + 不保留旧 Runtime,设置 `runtime_enabled=false`,也不持久化自由文本错误; + 下一轮 Snapshot Stream 据此继续发送完整 Snapshot。只有后续正向 ACK + 才把 `runtime_enabled` 重新设为 true,因此负向 ACK 后延迟到达的旧 Runtime + 不能重建已清除的报告与 sequence fence。 5. 旧版本、错误 checksum、超前版本或被替换 Session 均不修改状态。 `error_code` 只接受有界枚举风格字符串;`error_message` 只做长度校验,不写日志或 @@ -171,9 +194,10 @@ ACK 操作在 Redis Lua 的一个原子边界内完成: 处理顺序: 1. 校验 Session、报告序号、时间戳、计数非负、Proxy ID 唯一和批次上限。 -2. Snapshot version 与 ownership epoch 必须精确等于 Session 已 ACK 上界。 -3. 相同 sequence 和相同规范化内容是幂等重放;较小 sequence 拒绝;相同 sequence - 不同内容判定冲突。 +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` 为空。 @@ -211,6 +235,8 @@ controlPlane: - `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: