From 43dec7324a44eeb405b57e7de451ad1300a0f347 Mon Sep 17 00:00:00 2001 From: youfak Date: Fri, 31 Jul 2026 10:49:08 +0800 Subject: [PATCH] feat: add worker control plane configuration --- configs/proxy-pool.yaml | 3 + deploy/config/local.yaml | 3 + deploy/kubernetes/base/configmap.yaml | 2 + docs/configuration/reference.md | 76 ++++++++++++-- internal/config/config.go | 23 ++++ internal/config/config_test.go | 144 ++++++++++++++++++++++++++ internal/config/validate.go | 72 +++++++++++++ 7 files changed, 312 insertions(+), 11 deletions(-) diff --git a/configs/proxy-pool.yaml b/configs/proxy-pool.yaml index 5f35e1b..ba0ff50 100644 --- a/configs/proxy-pool.yaml +++ b/configs/proxy-pool.yaml @@ -81,6 +81,9 @@ admin: auth: mode: none +controlPlane: + enabled: false + metrics: enabled: true listen: 127.0.0.1:9090 diff --git a/deploy/config/local.yaml b/deploy/config/local.yaml index 17079df..fb30f37 100644 --- a/deploy/config/local.yaml +++ b/deploy/config/local.yaml @@ -58,6 +58,9 @@ admin: header: X-Admin-Key token: "${PROXY_POOL_ADMIN_TOKEN}" +controlPlane: + enabled: false + metrics: enabled: true listen: 0.0.0.0:9090 diff --git a/deploy/kubernetes/base/configmap.yaml b/deploy/kubernetes/base/configmap.yaml index 62116b9..5179992 100644 --- a/deploy/kubernetes/base/configmap.yaml +++ b/deploy/kubernetes/base/configmap.yaml @@ -60,6 +60,8 @@ data: mode: apiKey header: X-Admin-Key token: "${PROXY_POOL_ADMIN_TOKEN}" + controlPlane: + enabled: false metrics: enabled: true listen: 0.0.0.0:9090 diff --git a/docs/configuration/reference.md b/docs/configuration/reference.md index bf239f1..348d363 100644 --- a/docs/configuration/reference.md +++ b/docs/configuration/reference.md @@ -52,6 +52,7 @@ security: {} gateway: {} distribution: {} admin: {} +controlPlane: {} metrics: {} storage: {} routing: [] @@ -64,6 +65,7 @@ upstreams: {} - `gateway`:HTTP/HTTPS CONNECT 数据面入口。 - `distribution`:一次性独占提取入口。 - `admin`:运维管理入口,必须与 Distribution 分端口。 +- `controlPlane`:Worker 注册、Snapshot ACK 和 Runtime 心跳的 gRPC 控制面,默认关闭。 - `metrics`:Prometheus 入口。 - `storage`:Controller 使用的 PostgreSQL 与 Redis 地址。 - `routing`:有序 Routing 列表,自上而下首条命中停止。 @@ -122,7 +124,59 @@ Lua 可精确表示的整数范围内。 Gateway、Distribution、Admin 与 Provider API 是独立认证边界。改变其中一套 不得连带改变其他入口。 -## 4. Gateway +## 4. Worker 控制面 + +控制面默认关闭;默认配置中的 `8443` 端口预留不表示服务已监听。控制面 Session、 +Snapshot ACK 和 Runtime 报告是 Redis 的短效运行时状态,**不写入 PostgreSQL**。 + +仅本机开发可以使用回环明文: + +```yaml +controlPlane: + enabled: true + listen: 127.0.0.1:8443 + protocolVersion: 1 + heartbeatInterval: 10s + sessionTTL: 30s + maxStaleAge: 10s + maxMessageBytes: 1048576 + maxRuntimeCounters: 100000 + maxConcurrentStreams: 128 + tls: + mode: disabled +``` + +任何非回环监听地址都必须使用 mTLS: + +```yaml +controlPlane: + enabled: true + listen: 0.0.0.0:8443 + protocolVersion: 1 + heartbeatInterval: 10s + sessionTTL: 30s + maxStaleAge: 10s + maxMessageBytes: 1048576 + maxRuntimeCounters: 100000 + maxConcurrentStreams: 128 + tls: + mode: mtls + certFile: /run/secrets/controller-cert.pem + keyFile: /run/secrets/controller-key.pem + clientCAFile: /run/secrets/worker-ca.pem + trustDomain: proxy.example + environment: production +``` + +- `protocolVersion` 当前固定为 `1`。 +- `sessionTTL` 至少为 `3 * heartbeatInterval`;`maxStaleAge` 不得短于心跳间隔。 +- `maxMessageBytes` 范围为 `1..67108864`;`maxRuntimeCounters` 范围为 + `1..1000000`;`maxConcurrentStreams` 必须为正数。 +- `tls.mode` 只能为 `disabled` 或 `mtls`。`disabled` 只允许回环监听;`mtls` 必须 + 同时配置证书、私钥、客户端 CA、全小写 DNS `trustDomain` 和单 URI 路径段 + `environment`。 + +## 5. Gateway ```yaml gateway: @@ -151,7 +205,7 @@ gateway: - 保留地址、CGNAT 与云元数据端点始终拒绝,不能通过私网/链路本地开关放行。 - `maxConcurrentConnections` 是入口准入上限,不是 Proxy 容量上限。 -## 5. Distribution +## 6. Distribution ```yaml distribution: @@ -195,7 +249,7 @@ Extraction 是固定的一次性独占行为,**没有** `mode`、`leaseDuratio 启用认证。 - `authenticatedClientOrSourceIP`:优先认证主体,无主体时回退来源地址。 -## 6. Routing +## 7. Routing ```yaml routing: @@ -231,7 +285,7 @@ Sequential 的空计数属于 Upstream,当前索引属于 Routing。只有 Pro 成功、模板成功且合法候选为零时才增加空计数。错误不改变空计数;重复候选 会重置空计数但增加独立 duplicate 指标。 -## 7. Upstream +## 8. Upstream ```yaml upstreams: @@ -281,7 +335,7 @@ upstreams: urls: [http://connect.rom.miui.com/generate_204] ``` -### 7.1 Provider 与代理认证 +### 8.1 Provider 与代理认证 - `api.auth` 用于系统访问 Provider API。 - `proxyAuth` 用于最终连接被获取的 Proxy。 @@ -315,7 +369,7 @@ proxyAuth: password: "${PROVIDER_PROXY_PASSWORD}" ``` -### 7.2 Pool 与累计额度 +### 8.2 Pool 与累计额度 - `pool.maxSize`:当前系统维护且尚未 EXTRACTED 的 Proxy 硬上限,包括 FETCHED、CHECKING、AVAILABLE、SUSPECT、DRAINING 和 pending expected。 @@ -344,7 +398,7 @@ proxyAuth: 内存 lease 上限按所有配置 Upstream 的 `pool.maxSize * fetch.maxInFlight` 汇总, 配置 reload 只提高上限,不预分配对应内存。 -### 7.3 Fetch 限制 +### 8.3 Fetch 限制 - `estimatedIPsPerCall`:冷启动时每次 Provider 调用预计返回的合法 Proxy 数, 同时用于 `pool.maxSize` 的 pending 预占;不得从任意 Query 或 Body 字段推断。 @@ -359,7 +413,7 @@ proxyAuth: 大量缺池信号必须合并成 singleflight 或容量为 1 的通知,不能按 Gateway 请求 数量线性触发 Provider API。 -### 7.4 Refill 水位 +### 8.4 Refill 水位 - `reconcileInterval`:无事件时重新读取库存的兜底周期,启动时仍立即检查一次。 - `minimumAvailableSlots`:可用并发槽位低于该值时进入补池。 @@ -370,7 +424,7 @@ proxyAuth: `requestInterval` 只限制外部 API 调用,不能兼任库存复核周期。pending Proxy 按 `estimatedIPsPerCall * maxConcurrencyPerProxy` 折算槽位,避免并发补池超量。 -### 7.5 生命周期与健康 +### 8.5 生命周期与健康 - 明确绝对过期时间优先于响应 TTL,响应 TTL 优先于配置 `lifecycle.ttl`。 - 距离过期不足 `allocationSafetyMargin` 时停止新分配。 @@ -378,7 +432,7 @@ proxyAuth: - 第一次有意义失败进入 SUSPECT;达到 `maxConsecutiveFailures` 后才进入 UNHEALTHY。 -## 8. 存储、Admin 与 Metrics +## 9. 存储、Admin 与 Metrics ```yaml admin: @@ -415,7 +469,7 @@ Metrics 启用时 `listen` 必须是合法 `host:port`。该入口固定提供 ` 服务流量门槛,PostgreSQL 故障由 Admin 接口独立报告。Metrics 开关或监听地址 变更需要重启 Controller。 -## 9. 启动前校验清单 +## 10. 启动前校验清单 1. `version` 必须为 `1`,未知字段拒绝。 2. 所有启用监听器具有合法 `host:port`。 diff --git a/internal/config/config.go b/internal/config/config.go index 190065f..ff8c26a 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -31,6 +31,7 @@ type Config struct { 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"` @@ -127,6 +128,28 @@ type Metrics struct { Listen string `yaml:"listen"` } +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"` +} + type Storage struct { PostgresURL string `yaml:"postgresURL"` RedisURL string `yaml:"redisURL"` diff --git a/internal/config/config_test.go b/internal/config/config_test.go index ba2e55e..2a2b6a1 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -352,6 +352,150 @@ func TestValidateMetricsListener(t *testing.T) { } } +func TestValidateControlPlane(t *testing.T) { + tests := []struct { + name string + mutate func(*Config) + want string + }{ + { + name: "missing listen", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.Listen = "" + }, + want: "controlPlane listen", + }, + { + name: "public plaintext", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.Listen = "0.0.0.0:8443" + }, + want: "requires mtls", + }, + { + name: "short session ttl", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.SessionTTL = Duration(29 * time.Second) + }, + want: "sessionTTL", + }, + { + name: "stale below heartbeat", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.MaxStaleAge = Duration(9 * time.Second) + }, + want: "maxStaleAge", + }, + { + name: "invalid protocol", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.ProtocolVersion = 0 + }, + want: "protocolVersion", + }, + { + name: "invalid message limit", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.MaxMessageBytes = 64<<20 + 1 + }, + want: "maxMessageBytes", + }, + { + name: "invalid counter limit", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.MaxRuntimeCounters = MaximumPoolSize + 1 + }, + want: "maxRuntimeCounters", + }, + { + name: "zero concurrent streams", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.MaxConcurrentStreams = 0 + }, + want: "maxConcurrentStreams", + }, + { + name: "unsupported tls mode", + mutate: func(cfg *Config) { + cfg.ControlPlane = validControlPlane() + cfg.ControlPlane.TLS.Mode = "serverTLS" + }, + want: "tls.mode", + }, + { + name: "missing mtls files", + mutate: func(cfg *Config) { + cfg.ControlPlane = validMTLSControlPlane() + cfg.ControlPlane.TLS.CertFile = "" + }, + want: "certFile", + }, + { + name: "trust domain with port", + mutate: func(cfg *Config) { + cfg.ControlPlane = validMTLSControlPlane() + cfg.ControlPlane.TLS.TrustDomain = "proxy.example:443" + }, + want: "trustDomain", + }, + { + name: "environment path", + mutate: func(cfg *Config) { + cfg.ControlPlane = validMTLSControlPlane() + cfg.ControlPlane.TLS.Environment = "prod/eu" + }, + want: "environment", + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + cfg := mustLoadValidConfig(t) + test.mutate(cfg) + if err := Validate(cfg); err == nil || !strings.Contains(err.Error(), test.want) { + t.Fatalf("Validate() error = %v, want substring %q", err, test.want) + } + }) + } +} + +func validControlPlane() ControlPlane { + return ControlPlane{ + Enabled: true, + Listen: "127.0.0.1:8443", + ProtocolVersion: 1, + HeartbeatInterval: Duration(10 * time.Second), + SessionTTL: Duration(30 * time.Second), + MaxStaleAge: Duration(10 * time.Second), + MaxMessageBytes: 1 << 20, + MaxRuntimeCounters: 100_000, + MaxConcurrentStreams: 128, + TLS: ControlPlaneTLS{Mode: "disabled"}, + } +} + +func validMTLSControlPlane() ControlPlane { + controlPlane := validControlPlane() + controlPlane.Listen = "0.0.0.0:8443" + controlPlane.TLS = ControlPlaneTLS{ + Mode: "mtls", + CertFile: "/run/secrets/controller-cert.pem", + KeyFile: "/run/secrets/controller-key.pem", + ClientCAFile: "/run/secrets/worker-ca.pem", + TrustDomain: "proxy.example", + Environment: "production", + } + return controlPlane +} + func TestValidateRejectsInvalidConfigurationMatrix(t *testing.T) { tests := []struct { name string diff --git a/internal/config/validate.go b/internal/config/validate.go index b27d980..ac4b780 100644 --- a/internal/config/validate.go +++ b/internal/config/validate.go @@ -10,6 +10,11 @@ import ( "strings" ) +var ( + controlPlaneIdentityPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$`) + trustDomainLabelPattern = regexp.MustCompile(`^[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?$`) +) + func Validate(cfg *Config) error { if cfg == nil { return fmt.Errorf("validate configuration: nil config") @@ -38,6 +43,9 @@ func Validate(cfg *Config) error { return err } } + if err := validateControlPlane(cfg.ControlPlane); err != nil { + return err + } if fetchConfigured(cfg.Defaults.Fetch) { if err := validateFetch("defaults.fetch", cfg.Defaults.Fetch); err != nil { return err @@ -101,6 +109,70 @@ func Validate(cfg *Config) error { return nil } +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 := item.HeartbeatInterval.Value() + if heartbeat <= 0 { + return fmt.Errorf("validate controlPlane heartbeatInterval: must be greater than zero") + } + if item.SessionTTL.Value()/3 < heartbeat { + return fmt.Errorf("validate controlPlane sessionTTL: must be at least three heartbeat intervals") + } + if item.MaxStaleAge.Value() < 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 +} + +func validTrustDomain(value string) bool { + if len(value) == 0 || len(value) > 253 { + return false + } + for _, label := range strings.Split(value, ".") { + if !trustDomainLabelPattern.MatchString(label) { + return false + } + } + return true +} + func validateListener(name string, listener Listener, security Security) error { if !listener.Enabled { return nil