From 0d29b828bb13791d4d76cb1bf867f8a55ab50efc Mon Sep 17 00:00:00 2001 From: Rogee Date: Sat, 12 Sep 2026 11:09:42 +0800 Subject: [PATCH] feat: harden control-plane deployment --- control-plane/README.md | 29 ++- .../cmd/wxagent-control-plane/main.go | 82 +++++- .../cmd/wxagent-control-plane/main_test.go | 34 +++ control-plane/compose.production.yml | 56 ++++ control-plane/production_test.go | 197 +++++++++++++++ control-plane/server.go | 88 ++++++- control-plane/server_test.go | 6 + control-plane/store.go | 239 +++++++++++++++++- control-plane/tls.go | 94 +++++++ docs/WxAgent-远程协议草案-v1.0.md | 4 +- ...-远程多节点控制与白名单数据上报开发计划.md | 6 +- docs/WxAgent-远程控制面架构与部署说明-v1.0.md | 81 ++++-- node-agent/WxAgent.Core/RemoteContracts.cs | 61 ++++- .../WxAgent.Core/RemoteControlClient.cs | 50 +++- node-agent/WxAgent.Host/Program.cs | 2 + node-agent/WxAgent.Host/RemoteCliCommands.cs | 8 +- .../RemoteAgentHostedService.cs | 3 +- scripts/control-plane-data.sh | 42 +++ scripts/control-plane-scale-smoke.py | 124 +++++++++ scripts/remote-control-smoke.sh | 2 +- scripts/revoke-client-certificate.sh | 30 +++ .../WxAgent.Core.Tests/RemoteControlTests.cs | 46 ++++ 22 files changed, 1213 insertions(+), 71 deletions(-) create mode 100644 control-plane/cmd/wxagent-control-plane/main_test.go create mode 100644 control-plane/compose.production.yml create mode 100644 control-plane/production_test.go create mode 100644 control-plane/tls.go create mode 100755 scripts/control-plane-data.sh create mode 100755 scripts/control-plane-scale-smoke.py create mode 100755 scripts/revoke-client-certificate.sh diff --git a/control-plane/README.md b/control-plane/README.md index 34d5824..df72c2c 100644 --- a/control-plane/README.md +++ b/control-plane/README.md @@ -27,9 +27,23 @@ export WXAGENT_NODE_TOKENS='{"node-a":"token-a","node-b":"token-b"}' export WXAGENT_WEB_USERS='{"admin":"change-me","auditor":"read-only-password"}' ``` -可选:`WXAGENT_CONTROL_PLANE_ADDR`(默认 `127.0.0.1:8090`)、`WXAGENT_CONTROL_PLANE_DATA`(默认 `control-plane-data.json`)。Token 只从环境变量读取,不写入控制面数据文件和日志。节点 `service.json` 可设置 `remoteConfigurationFile` 指向 CLI 管理的 `remote.json`;修改上报配置会在下一轮生效,修改远程地址/Token 后需重启 Agent。 +可选:`WXAGENT_CONTROL_PLANE_ADDR`(默认 `127.0.0.1:8090`)、`WXAGENT_CONTROL_PLANE_DATA`(默认 `control-plane-data.json`)。凭据可通过环境变量或 `*_FILE` Secret 文件读取,不写入控制面数据文件和日志。节点 `service.json` 可设置 `remoteConfigurationFile` 指向 CLI 管理的 `remote.json`;修改上报配置会在下一轮生效,修改远程地址/Token 后需重启 Agent。 -浏览器访问 `http://127.0.0.1:8090/`,登录后管理节点、任务、白名单事件和审计记录。生产部署必须使用 HTTPS 和外部密钥管理;本地 HTTP 仅用于 loopback 集成测试。 +生产控制面可直接启用 TLS 和节点 mTLS: + +```text +WXAGENT_CONTROL_PLANE_TLS_CERT_FILE +WXAGENT_CONTROL_PLANE_TLS_KEY_FILE +WXAGENT_MTLS_CLIENT_CA_FILE +WXAGENT_MTLS_REQUIRE_NODE_CERT=true +WXAGENT_MTLS_REVOKED_CERTS_FILE +WXAGENT_NODE_TOKENS_FILE=/run/secrets/node_tokens +WXAGENT_WEB_USERS_FILE=/run/secrets/web_users +``` + +节点证书按证书原始 DER 的 SHA-256 指纹逐行写入撤销文件;文件读取失败时节点认证拒绝。控制面还提供 `/readyz`,并对数据文件使用单写入者锁、自动轮转备份和任务/事件/审计保留期限。生产示例见 [`compose.production.yml`](./compose.production.yml)。 + +浏览器访问 `http://127.0.0.1:8090/`,登录后管理节点、任务、白名单事件和审计记录。生产部署必须使用 HTTPS、节点 mTLS 和外部 Secret 文件;本地 HTTP 仅用于显式的私有网络/loopback 集成测试。 远程只读查询使用同一持久化任务队列,返回 `task_id` 后通过 `GET /v1/tasks/{task_id}` 取结果: @@ -55,6 +69,15 @@ cd ../.. `remote-control-smoke.sh` 会启动一次临时控制面,运行 .NET 节点协议客户端,验证注册、心跳、任务租约、幂等、结果回传和白名单拒绝;不会操作真实微信联系人。 +控制面并发/读取任务基线(使用已注册的测试节点,不执行写任务): + +```bash +WXAGENT_SCALE_WEB_PASSWORD='' \ + scripts/control-plane-scale-smoke.py --base-url http://127.0.0.1:8090 \ + --node-id node-a --account-id account-a --requests 100 --workers 10 +``` + +`control-plane-data.sh backup`/`restore` 用于停机前后的手工备份和原子恢复;自动备份由 Store 按 `WXAGENT_CONTROL_PLANE_BACKUP_INTERVAL`(默认 5 分钟)产生。当前 JSON 存储使用 active/passive 单写锁,不支持 active-active 多实例共享写入。 ## Windows 手工验收 1. 在已登录且未锁定的微信桌面会话中运行 `WxAgent.Host doctor` 和 `WxAgent.Host inspect-ui --output artifacts/ui-tree.json`,确认账号绑定使用稳定 `accountId`,会话使用稳定 `AutomationId`,不使用昵称/PID/窗口句柄猜测。 @@ -78,3 +101,5 @@ cd ../.. ## 容器镜像 `.gitea/workflows/build-web-image.yml` 会构建并推送 `git.ipao.vip//` 镜像:主分支额外更新 `latest`,`v*` 标签额外更新对应版本标签。请在 Gitea 账号级 Secrets 配置 `REGISTRY_TOKEN`;登录用户名使用仓库所有者。 + +生产环境不要提交 `compose.production.yml` 所引用的 `secrets/` 文件。将节点 Token map、Web 用户 map、TLS 私钥、客户端 CA 和撤销指纹列表放入外部 Secret 管理或受 ACL 保护的部署目录,再运行 `docker compose -f compose.production.yml up -d`。证书泄露时可用 `scripts/revoke-client-certificate.sh` 更新指纹列表,并重建 Docker Secret 挂载的容器。 diff --git a/control-plane/cmd/wxagent-control-plane/main.go b/control-plane/cmd/wxagent-control-plane/main.go index 56c75bd..8b8dfc7 100644 --- a/control-plane/cmd/wxagent-control-plane/main.go +++ b/control-plane/cmd/wxagent-control-plane/main.go @@ -3,12 +3,13 @@ package main import ( "context" "encoding/json" - "fmt" "log" "os" "os/signal" + "strconv" "strings" "syscall" + "time" controlplane "git.ipao.vip/rogee/wx-win-agent/control-plane" ) @@ -21,22 +22,34 @@ func main() { if value := os.Getenv("WXAGENT_CONTROL_PLANE_DATA"); value != "" { config.DataFile = value } - config.NodeTokens = readMapEnv("WXAGENT_NODE_TOKENS") + config.BackupDir = os.Getenv("WXAGENT_CONTROL_PLANE_BACKUP_DIR") + config.BackupCount = readIntEnv("WXAGENT_CONTROL_PLANE_BACKUP_COUNT", config.BackupCount) + config.BackupInterval = readDurationEnv("WXAGENT_CONTROL_PLANE_BACKUP_INTERVAL", config.BackupInterval) + config.TaskRetention = readDurationEnv("WXAGENT_CONTROL_PLANE_TASK_RETENTION", config.TaskRetention) + config.EventRetention = readDurationEnv("WXAGENT_CONTROL_PLANE_EVENT_RETENTION", config.EventRetention) + config.AuditRetention = readDurationEnv("WXAGENT_CONTROL_PLANE_AUDIT_RETENTION", config.AuditRetention) + config.TLSCertFile = os.Getenv("WXAGENT_CONTROL_PLANE_TLS_CERT_FILE") + config.TLSKeyFile = os.Getenv("WXAGENT_CONTROL_PLANE_TLS_KEY_FILE") + config.MTLSClientCAFile = os.Getenv("WXAGENT_MTLS_CLIENT_CA_FILE") + config.MTLSRequireNodeCert = readBoolEnv("WXAGENT_MTLS_REQUIRE_NODE_CERT", false) + config.MTLSRevokedCertsFile = os.Getenv("WXAGENT_MTLS_REVOKED_CERTS_FILE") + + config.NodeTokens = readMapSecret("WXAGENT_NODE_TOKENS", "WXAGENT_NODE_TOKENS_FILE") if len(config.NodeTokens) == 0 { - nodeID, nodeToken := os.Getenv("WXAGENT_NODE_ID"), os.Getenv("WXAGENT_NODE_TOKEN") + nodeID, nodeToken := os.Getenv("WXAGENT_NODE_ID"), readSecret("WXAGENT_NODE_TOKEN", "WXAGENT_NODE_TOKEN_FILE") if strings.TrimSpace(nodeID) != "" && nodeToken != "" { config.NodeTokens = map[string]string{nodeID: nodeToken} } } - config.WebUsers = readMapEnv("WXAGENT_WEB_USERS") + config.WebUsers = readMapSecret("WXAGENT_WEB_USERS", "WXAGENT_WEB_USERS_FILE") if len(config.WebUsers) == 0 { - webUser, webPassword := os.Getenv("WXAGENT_WEB_USER"), os.Getenv("WXAGENT_WEB_PASSWORD") + webUser, webPassword := os.Getenv("WXAGENT_WEB_USER"), readSecret("WXAGENT_WEB_PASSWORD", "WXAGENT_WEB_PASSWORD_FILE") if strings.TrimSpace(webUser) != "" && webPassword != "" { config.WebUsers = map[string]string{webUser: webPassword} } } if len(config.NodeTokens) == 0 || len(config.WebUsers) == 0 { - log.Fatal("configure WXAGENT_NODE_TOKENS/WXAGENT_WEB_USERS as JSON maps, or the single-node fallback environment variables") + log.Fatal("configure node and web credentials with environment variables or *_FILE secret files") } server, err := controlplane.NewServer(config) @@ -47,19 +60,66 @@ func main() { defer stop() log.Printf("wxagent control plane listening on %s; data=%s", config.ListenAddr, config.DataFile) if err := server.ListenAndServe(ctx); err != nil { - fmt.Fprintln(os.Stderr, err) + log.Printf("control plane stopped: %v", err) os.Exit(1) } } -func readMapEnv(name string) map[string]string { - value := strings.TrimSpace(os.Getenv(name)) +func readSecret(valueName, fileName string) string { + if path := strings.TrimSpace(os.Getenv(fileName)); path != "" { + data, err := os.ReadFile(path) + if err != nil { + log.Fatalf("read %s: %v", fileName, err) + } + return strings.TrimRight(string(data), "\r\n") + } + return os.Getenv(valueName) +} + +func readMapSecret(valueName, fileName string) map[string]string { + value := readSecret(valueName, fileName) if value == "" { return nil } - var result map[string]string + result := make(map[string]string) if err := json.Unmarshal([]byte(value), &result); err != nil { - log.Fatalf("%s must be a JSON object", name) + log.Fatalf("%s must contain a JSON object", fileName) } return result } + +func readBoolEnv(name string, fallback bool) bool { + value := strings.TrimSpace(os.Getenv(name)) + if value == "" { + return fallback + } + parsed, err := strconv.ParseBool(value) + if err != nil { + log.Fatalf("%s must be true or false", name) + } + return parsed +} + +func readIntEnv(name string, fallback int) int { + value := strings.TrimSpace(os.Getenv(name)) + if value == "" { + return fallback + } + parsed, err := strconv.Atoi(value) + if err != nil || parsed < 1 { + log.Fatalf("%s must be a positive integer", name) + } + return parsed +} + +func readDurationEnv(name string, fallback time.Duration) time.Duration { + value := strings.TrimSpace(os.Getenv(name)) + if value == "" { + return fallback + } + parsed, err := time.ParseDuration(value) + if err != nil || parsed < 0 { + log.Fatalf("%s must be a non-negative duration", name) + } + return parsed +} diff --git a/control-plane/cmd/wxagent-control-plane/main_test.go b/control-plane/cmd/wxagent-control-plane/main_test.go new file mode 100644 index 0000000..62f23d8 --- /dev/null +++ b/control-plane/cmd/wxagent-control-plane/main_test.go @@ -0,0 +1,34 @@ +package main + +import ( + "os" + "path/filepath" + "testing" +) + +func TestReadSecretPrefersFileAndRemovesOnlyLineEndings(t *testing.T) { + directory := t.TempDir() + path := filepath.Join(directory, "token") + if err := os.WriteFile(path, []byte("file-token\r\n"), 0o600); err != nil { + t.Fatal(err) + } + t.Setenv("TEST_TOKEN", "environment-token") + t.Setenv("TEST_TOKEN_FILE", path) + if got := readSecret("TEST_TOKEN", "TEST_TOKEN_FILE"); got != "file-token" { + t.Fatalf("readSecret() = %q", got) + } +} + +func TestReadMapSecretFromFile(t *testing.T) { + directory := t.TempDir() + path := filepath.Join(directory, "users.json") + if err := os.WriteFile(path, []byte("{\"admin\":\"password\"}\n"), 0o600); err != nil { + t.Fatal(err) + } + t.Setenv("TEST_USERS", "") + t.Setenv("TEST_USERS_FILE", path) + users := readMapSecret("TEST_USERS", "TEST_USERS_FILE") + if users["admin"] != "password" || len(users) != 1 { + t.Fatalf("readMapSecret() = %#v", users) + } +} diff --git a/control-plane/compose.production.yml b/control-plane/compose.production.yml new file mode 100644 index 0000000..1618d7a --- /dev/null +++ b/control-plane/compose.production.yml @@ -0,0 +1,56 @@ +services: + control-plane: + image: ${WXAGENT_CONTROL_PLANE_IMAGE:?set image tag} + restart: unless-stopped + read_only: true + user: "65532:65532" + security_opt: + - no-new-privileges:true + cap_drop: + - ALL + tmpfs: + - /tmp:rw,noexec,nosuid,size=16m + ports: + - "${CONTROL_PLANE_BIND:-127.0.0.1:18090}:8090" + volumes: + - control-plane-data:/data + secrets: + - node_tokens + - web_users + - tls_cert + - tls_key + - client_ca + - revoked_client_certificates + environment: + WXAGENT_CONTROL_PLANE_ADDR: 0.0.0.0:8090 + WXAGENT_CONTROL_PLANE_DATA: /data/control-plane-data.json + WXAGENT_CONTROL_PLANE_BACKUP_DIR: /data/backups + WXAGENT_CONTROL_PLANE_BACKUP_COUNT: "7" + WXAGENT_CONTROL_PLANE_BACKUP_INTERVAL: 5m + WXAGENT_CONTROL_PLANE_TASK_RETENTION: 720h + WXAGENT_CONTROL_PLANE_EVENT_RETENTION: 720h + WXAGENT_CONTROL_PLANE_AUDIT_RETENTION: 2160h + WXAGENT_CONTROL_PLANE_TLS_CERT_FILE: /run/secrets/tls_cert + WXAGENT_CONTROL_PLANE_TLS_KEY_FILE: /run/secrets/tls_key + WXAGENT_MTLS_CLIENT_CA_FILE: /run/secrets/client_ca + WXAGENT_MTLS_REQUIRE_NODE_CERT: "true" + WXAGENT_MTLS_REVOKED_CERTS_FILE: /run/secrets/revoked_client_certificates + WXAGENT_NODE_TOKENS_FILE: /run/secrets/node_tokens + WXAGENT_WEB_USERS_FILE: /run/secrets/web_users + +secrets: + node_tokens: + file: ./secrets/node-tokens.json + web_users: + file: ./secrets/web-users.json + tls_cert: + file: ./secrets/server.crt + tls_key: + file: ./secrets/server.key + client_ca: + file: ./secrets/client-ca.crt + revoked_client_certificates: + file: ./secrets/revoked-client-certificates.txt + +volumes: + control-plane-data: diff --git a/control-plane/production_test.go b/control-plane/production_test.go new file mode 100644 index 0000000..0779723 --- /dev/null +++ b/control-plane/production_test.go @@ -0,0 +1,197 @@ +package controlplane + +import ( + "crypto/rand" + "crypto/rsa" + "crypto/sha256" + "crypto/tls" + "crypto/x509" + "crypto/x509/pkix" + "encoding/hex" + "encoding/pem" + "math/big" + "net/http" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +func TestStoreLocksDataAndKeepsBoundedBackups(t *testing.T) { + directory := t.TempDir() + path := filepath.Join(directory, "state.json") + store, err := OpenStore(path, StoreOptions{BackupDir: filepath.Join(directory, "backups"), BackupCount: 2}) + if err != nil { + t.Fatal(err) + } + if err := store.Mutate(func(state *PersistedState) error { + state.Nodes["node-1"] = Node{NodeID: "node-1"} + return nil + }); err != nil { + t.Fatal(err) + } + if err := store.Mutate(func(state *PersistedState) error { + state.Nodes["node-1"] = Node{NodeID: "node-1", AgentVersion: "2"} + return nil + }); err != nil { + t.Fatal(err) + } + if err := store.Mutate(func(state *PersistedState) error { + state.Nodes["node-1"] = Node{NodeID: "node-1", AgentVersion: "3"} + return nil + }); err != nil { + t.Fatal(err) + } + if _, err := OpenStore(path, StoreOptions{BackupDir: filepath.Join(directory, "backups"), BackupCount: 2}); err == nil || !strings.Contains(err.Error(), "already in use") { + t.Fatalf("second writer was not rejected: %v", err) + } + if err := store.Close(); err != nil { + t.Fatal(err) + } + entries, err := os.ReadDir(filepath.Join(directory, "backups")) + if err != nil { + t.Fatal(err) + } + if len(entries) != 2 { + t.Fatalf("expected two backups, got %d", len(entries)) + } + for _, entry := range entries { + if entry.Type().Perm() != 0o600 { + info, infoErr := entry.Info() + if infoErr != nil { + t.Fatal(infoErr) + } + if info.Mode().Perm() != 0o600 { + t.Fatalf("backup permissions = %o", info.Mode().Perm()) + } + } + } + second, err := OpenStore(path, StoreOptions{BackupDir: filepath.Join(directory, "backups"), BackupCount: 2}) + if err != nil { + t.Fatal(err) + } + defer second.Close() +} + +func TestStorePrunesOnlyTerminalRecordsByRetention(t *testing.T) { + directory := t.TempDir() + path := filepath.Join(directory, "state.json") + seed, err := OpenStore(path, StoreOptions{BackupDir: "", BackupCount: -1}) + if err != nil { + t.Fatal(err) + } + old := time.Now().UTC().Add(-2 * time.Hour) + if err := seed.Mutate(func(state *PersistedState) error { + state.Tasks["old-terminal"] = Task{TaskID: "old-terminal", Status: TaskSucceeded, UpdatedAt: old} + state.Tasks["old-running"] = Task{TaskID: "old-running", Status: TaskRunning, UpdatedAt: old} + state.Events = append(state.Events, StoredEvent{ReceivedAt: old}) + state.Audit = append(state.Audit, AuditEntry{At: old}) + return nil + }); err != nil { + t.Fatal(err) + } + if err := seed.Close(); err != nil { + t.Fatal(err) + } + + store, err := OpenStore(path, StoreOptions{ + BackupDir: filepath.Join(directory, "backups"), BackupCount: -1, + TaskRetention: time.Hour, EventRetention: time.Hour, AuditRetention: time.Hour, + }) + if err != nil { + t.Fatal(err) + } + defer store.Close() + if err := store.Mutate(func(state *PersistedState) error { + state.Nodes["node-1"] = Node{NodeID: "node-1"} + return nil + }); err != nil { + t.Fatal(err) + } + snapshot := store.Snapshot() + if _, ok := snapshot.Tasks["old-terminal"]; ok { + t.Fatal("old terminal task was retained") + } + if _, ok := snapshot.Tasks["old-running"]; !ok { + t.Fatal("old running task was pruned") + } + if len(snapshot.Events) != 0 || len(snapshot.Audit) != 0 { + t.Fatalf("old events/audit were retained: events=%d audit=%d", len(snapshot.Events), len(snapshot.Audit)) + } +} + +func TestTLSConfigSupportsVerifiedNodeClientCertificates(t *testing.T) { + directory := t.TempDir() + certFile, keyFile, caFile := writeTestCertificateMaterial(t, directory) + config := ServerConfig{ + TLSCertFile: certFile, + TLSKeyFile: keyFile, + MTLSClientCAFile: caFile, + MTLSRequireNodeCert: true, + MTLSRevokedCertsFile: filepath.Join(directory, "revoked.txt"), + } + if err := os.WriteFile(config.MTLSRevokedCertsFile, []byte("# initially empty\n"), 0o600); err != nil { + t.Fatal(err) + } + tlsConfig, err := newTLSConfig(config) + if err != nil { + t.Fatal(err) + } + if tlsConfig.MinVersion != tls.VersionTLS13 || tlsConfig.ClientAuth != tls.VerifyClientCertIfGiven || len(tlsConfig.Certificates) != 1 { + t.Fatalf("unexpected TLS config: min=%d auth=%d certs=%d", tlsConfig.MinVersion, tlsConfig.ClientAuth, len(tlsConfig.Certificates)) + } + + rawCertificate := []byte("client-cert") + digest := sha256.Sum256(rawCertificate) + fingerprint := hex.EncodeToString(digest[:]) + if err := os.WriteFile(config.MTLSRevokedCertsFile, []byte(fingerprint+"\n"), 0o600); err != nil { + t.Fatal(err) + } + server := &Server{config: config} + request := &http.Request{TLS: &tls.ConnectionState{PeerCertificates: []*x509.Certificate{{Raw: rawCertificate}}}} + if server.clientCertificateAllowed(request) { + t.Fatal("revoked certificate was accepted") + } +} + +func TestTLSConfigRejectsIncompleteSecuritySettings(t *testing.T) { + if _, err := newTLSConfig(ServerConfig{TLSCertFile: "cert.pem"}); err == nil { + t.Fatal("incomplete TLS certificate settings were accepted") + } + if _, err := newTLSConfig(ServerConfig{MTLSRequireNodeCert: true}); err == nil { + t.Fatal("mTLS without TLS material was accepted") + } +} + +func writeTestCertificateMaterial(t *testing.T, directory string) (string, string, string) { + t.Helper() + key, err := rsa.GenerateKey(rand.Reader, 2048) + if err != nil { + t.Fatal(err) + } + template := &x509.Certificate{ + SerialNumber: big.NewInt(1), + Subject: pkix.Name{CommonName: "wxagent-test"}, + NotBefore: time.Now().Add(-time.Minute), + NotAfter: time.Now().Add(time.Hour), + KeyUsage: x509.KeyUsageDigitalSignature | x509.KeyUsageKeyEncipherment | x509.KeyUsageCertSign, + IsCA: true, + BasicConstraintsValid: true, + } + der, err := x509.CreateCertificate(rand.Reader, template, template, &key.PublicKey, key) + if err != nil { + t.Fatal(err) + } + certFile := filepath.Join(directory, "server.pem") + keyFile := filepath.Join(directory, "server-key.pem") + caFile := filepath.Join(directory, "client-ca.pem") + certPEM := pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}) + keyPEM := pem.EncodeToMemory(&pem.Block{Type: "RSA PRIVATE KEY", Bytes: x509.MarshalPKCS1PrivateKey(key)}) + for path, data := range map[string][]byte{certFile: certPEM, keyFile: keyPEM, caFile: certPEM} { + if err := os.WriteFile(path, data, 0o600); err != nil { + t.Fatal(err) + } + } + return certFile, keyFile, caFile +} diff --git a/control-plane/server.go b/control-plane/server.go index abcf2da..514bd7f 100644 --- a/control-plane/server.go +++ b/control-plane/server.go @@ -5,6 +5,7 @@ import ( "crypto/rand" "crypto/sha256" "crypto/subtle" + "crypto/tls" "encoding/base64" "encoding/hex" "encoding/json" @@ -26,13 +27,24 @@ import ( ) type ServerConfig struct { - ListenAddr string - DataFile string - NodeTokens map[string]string - WebUsers map[string]string - LeaseTTL time.Duration - HeartbeatTimeout time.Duration - SessionTTL time.Duration + ListenAddr string + DataFile string + NodeTokens map[string]string + WebUsers map[string]string + LeaseTTL time.Duration + HeartbeatTimeout time.Duration + SessionTTL time.Duration + TLSCertFile string + TLSKeyFile string + MTLSClientCAFile string + MTLSRequireNodeCert bool + MTLSRevokedCertsFile string + BackupDir string + BackupCount int + BackupInterval time.Duration + TaskRetention time.Duration + EventRetention time.Duration + AuditRetention time.Duration } func DefaultServerConfig() ServerConfig { @@ -44,12 +56,18 @@ func DefaultServerConfig() ServerConfig { LeaseTTL: 30 * time.Second, HeartbeatTimeout: 45 * time.Second, SessionTTL: 8 * time.Hour, + BackupCount: 7, + BackupInterval: 5 * time.Minute, + TaskRetention: 30 * 24 * time.Hour, + EventRetention: 30 * 24 * time.Hour, + AuditRetention: 90 * 24 * time.Hour, } } type Server struct { config ServerConfig store *Store + tlsConfig *tls.Config nodeTokenHashes map[string][32]byte userPasswords map[string][32]byte sessionMu sync.Mutex @@ -77,6 +95,24 @@ func NewServer(config ServerConfig) (*Server, error) { if config.DataFile == "" { config.DataFile = defaults.DataFile } + if config.BackupDir == "" { + config.BackupDir = config.DataFile + ".backups" + } + if config.BackupCount == 0 { + config.BackupCount = defaults.BackupCount + } + if config.BackupInterval == 0 { + config.BackupInterval = defaults.BackupInterval + } + if config.TaskRetention == 0 { + config.TaskRetention = defaults.TaskRetention + } + if config.EventRetention == 0 { + config.EventRetention = defaults.EventRetention + } + if config.AuditRetention == 0 { + config.AuditRetention = defaults.AuditRetention + } if config.LeaseTTL <= 0 { config.LeaseTTL = defaults.LeaseTTL } @@ -92,13 +128,24 @@ func NewServer(config ServerConfig) (*Server, error) { if config.WebUsers == nil { config.WebUsers = map[string]string{} } - store, err := OpenStore(config.DataFile) + if config.BackupCount < 0 || config.BackupInterval < 0 || config.TaskRetention < 0 || config.EventRetention < 0 || config.AuditRetention < 0 { + return nil, errors.New("backup count and retention settings must be non-negative") + } + tlsConfig, err := newTLSConfig(config) + if err != nil { + return nil, err + } + store, err := OpenStore(config.DataFile, StoreOptions{ + BackupDir: config.BackupDir, BackupCount: config.BackupCount, BackupInterval: config.BackupInterval, + TaskRetention: config.TaskRetention, EventRetention: config.EventRetention, AuditRetention: config.AuditRetention, + }) if err != nil { return nil, err } s := &Server{ config: config, store: store, + tlsConfig: tlsConfig, nodeTokenHashes: map[string][32]byte{}, userPasswords: map[string][32]byte{}, sessions: map[string]session{}, @@ -126,13 +173,23 @@ func (s *Server) ListenAndServe(ctx context.Context) error { defer cancel() _ = server.Shutdown(shutdownCtx) }() - err := server.ListenAndServe() + defer func() { _ = s.Close() }() + var err error + if s.tlsConfig != nil { + server.TLSConfig = s.tlsConfig + err = server.ListenAndServeTLS("", "") + } else { + err = server.ListenAndServe() + } if errors.Is(err, http.ErrServerClosed) { return nil } return err } +// Close releases the active/passive store lock. +func (s *Server) Close() error { return s.store.Close() } + func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { correlationID := r.Header.Get("X-Correlation-Id") if !validIdentifier(correlationID, 128) { @@ -158,6 +215,8 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { case r.URL.Path == "/healthz" && r.Method == http.MethodGet: writeJSON(w, http.StatusOK, map[string]any{"status": "ok", "protocol_version": ProtocolVersion, "correlation_id": correlationID}) return + case r.URL.Path == "/readyz" && r.Method == http.MethodGet: + err = s.ready(w, correlationID) case r.URL.Path == "/v1/auth/login" && r.Method == http.MethodPost: err = s.login(w, r, correlationID) case r.URL.Path == "/v1/nodes" && r.Method == http.MethodGet: @@ -182,6 +241,14 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { } } +func (s *Server) ready(w http.ResponseWriter, correlationID string) error { + if err := s.store.Read(func(PersistedState) error { return nil }); err != nil { + return requestError{status: http.StatusServiceUnavailable, code: "NotReady", message: "The control plane store is not ready."} + } + writeJSON(w, http.StatusOK, map[string]any{"status": "ready", "protocol_version": ProtocolVersion, "correlation_id": correlationID}) + return nil +} + func (s *Server) login(w http.ResponseWriter, r *http.Request, correlationID string) error { var request struct { Username string `json:"username"` @@ -1071,6 +1138,9 @@ func (s *Server) auditRoute(w http.ResponseWriter, r *http.Request, _ string) er } func (s *Server) authenticateNode(r *http.Request) (string, error) { + if !s.clientCertificateAllowed(r) { + return "", requestError{status: http.StatusUnauthorized, code: "Unauthorized", message: "Node certificate authentication failed."} + } token := bearerToken(r) if token == "" { return "", requestError{status: http.StatusUnauthorized, code: "Unauthorized", message: "Node authentication is required."} diff --git a/control-plane/server_test.go b/control-plane/server_test.go index c8b4221..4adcb85 100644 --- a/control-plane/server_test.go +++ b/control-plane/server_test.go @@ -19,6 +19,7 @@ func TestReactFrontendIsEmbedded(t *testing.T) { } httpServer := httptest.NewServer(server.Handler()) defer httpServer.Close() + defer server.Close() response, err := http.Get(httpServer.URL + "/") if err != nil { t.Fatal(err) @@ -70,6 +71,7 @@ func TestNodeWebTaskAndEventFlow(t *testing.T) { } httpServer := httptest.NewServer(server.Handler()) defer httpServer.Close() + defer server.Close() client := httpServer.Client() if response := doJSON(t, client, http.MethodGet, httpServer.URL+"/v1/nodes", "", nil); response.Code != http.StatusUnauthorized { @@ -195,6 +197,9 @@ func TestNodeWebTaskAndEventFlow(t *testing.T) { t.Fatalf("expected two account-scoped events, got %d", len(eventList.Events)) } + if err := server.Close(); err != nil { + t.Fatal(err) + } restarted, err := NewServer(ServerConfig{DataFile: dataFile, NodeTokens: map[string]string{"node-1": "node-secret"}, WebUsers: map[string]string{"admin": "web-secret"}}) if err != nil { t.Fatal(err) @@ -226,6 +231,7 @@ func TestRemoteReadRoutesUseTheTaskQueueAndRetainReadContent(t *testing.T) { } httpServer := httptest.NewServer(server.Handler()) defer httpServer.Close() + defer server.Close() client := httpServer.Client() login := doJSON(t, client, http.MethodPost, httpServer.URL+"/v1/auth/login", "", map[string]string{"username": "admin", "password": "web-secret"}) diff --git a/control-plane/store.go b/control-plane/store.go index 2b9b1d8..6bd81ea 100644 --- a/control-plane/store.go +++ b/control-plane/store.go @@ -6,33 +6,88 @@ import ( "fmt" "os" "path/filepath" + "sort" + "strings" "sync" + "syscall" + "time" ) -// Store is a small durable JSON store for the single-process control-plane MVP. -// The file is an implementation detail; callers only observe transactional methods. -type Store struct { - path string - mu sync.Mutex - state PersistedState +// StoreOptions controls durability and bounded retention for the JSON store. +// A Store still has one active writer; the lock file enables active/passive failover. +type StoreOptions struct { + BackupDir string + BackupCount int + BackupInterval time.Duration + TaskRetention time.Duration + EventRetention time.Duration + AuditRetention time.Duration } -func OpenStore(path string) (*Store, error) { +// Store is a small durable JSON store for the control-plane MVP. +// The file is an implementation detail; callers only observe transactional methods. +type Store struct { + path string + lockFile *os.File + options StoreOptions + lastBackup time.Time + mu sync.Mutex + state PersistedState + closed bool +} + +func OpenStore(path string, options ...StoreOptions) (*Store, error) { if path == "" { return nil, errors.New("data file is required") } - store := &Store{path: filepath.Clean(path), state: PersistedState{Nodes: map[string]Node{}, Tasks: map[string]Task{}, Events: []StoredEvent{}, Audit: []AuditEntry{}}} + cleanPath := filepath.Clean(path) + storeOptions := StoreOptions{} + if len(options) > 0 { + storeOptions = options[0] + } + if storeOptions.BackupDir == "" { + storeOptions.BackupDir = cleanPath + ".backups" + } + if storeOptions.BackupCount == 0 { + storeOptions.BackupCount = 7 + } + if storeOptions.BackupCount < -1 { + return nil, errors.New("backup count must be -1 or non-negative") + } + if err := os.MkdirAll(filepath.Dir(cleanPath), 0o700); err != nil { + return nil, fmt.Errorf("create data directory: %w", err) + } + lockFile, err := os.OpenFile(cleanPath+".lock", os.O_CREATE|os.O_RDWR, 0o600) + if err != nil { + return nil, fmt.Errorf("open data lock: %w", err) + } + if err := syscall.Flock(int(lockFile.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); err != nil { + _ = lockFile.Close() + if errors.Is(err, syscall.EWOULDBLOCK) { + return nil, fmt.Errorf("data file is already in use: %s", cleanPath) + } + return nil, fmt.Errorf("lock data file: %w", err) + } + + store := &Store{ + path: cleanPath, + lockFile: lockFile, + options: storeOptions, + state: PersistedState{Nodes: map[string]Node{}, Tasks: map[string]Task{}, Events: []StoredEvent{}, Audit: []AuditEntry{}}, + } data, err := os.ReadFile(store.path) if errors.Is(err, os.ErrNotExist) { return store, nil } if err != nil { + _ = store.Close() return nil, fmt.Errorf("read data file: %w", err) } if len(data) == 0 { return store, nil } if err := json.Unmarshal(data, &store.state); err != nil { + _ = store.Close() return nil, fmt.Errorf("decode data file: %w", err) } store.ensureMaps() @@ -48,6 +103,9 @@ func (s *Store) Snapshot() PersistedState { func (s *Store) Mutate(fn func(*PersistedState) error) error { s.mu.Lock() defer s.mu.Unlock() + if s.closed { + return errors.New("store is closed") + } s.ensureMaps() if err := fn(&s.state); err != nil { return err @@ -58,9 +116,35 @@ func (s *Store) Mutate(fn func(*PersistedState) error) error { func (s *Store) Read(fn func(PersistedState) error) error { s.mu.Lock() defer s.mu.Unlock() + if s.closed { + return errors.New("store is closed") + } return fn(cloneState(s.state)) } +// Close releases the single-writer lock. A second control-plane process can then +// be promoted by the supervisor using the same data directory. +func (s *Store) Close() error { + s.mu.Lock() + defer s.mu.Unlock() + if s.closed { + return nil + } + s.closed = true + if s.lockFile == nil { + return nil + } + unlockErr := syscall.Flock(int(s.lockFile.Fd()), syscall.LOCK_UN) + closeErr := s.lockFile.Close() + if unlockErr != nil { + return fmt.Errorf("unlock data file: %w", unlockErr) + } + if closeErr != nil { + return fmt.Errorf("close data lock: %w", closeErr) + } + return nil +} + func (s *Store) ensureMaps() { if s.state.Nodes == nil { s.state.Nodes = map[string]Node{} @@ -80,6 +164,10 @@ func (s *Store) saveLocked() error { if err := os.MkdirAll(filepath.Dir(s.path), 0o700); err != nil { return fmt.Errorf("create data directory: %w", err) } + s.pruneLocked(time.Now().UTC()) + if err := s.backupLocked(); err != nil { + return err + } data, err := json.MarshalIndent(s.state, "", " ") if err != nil { return fmt.Errorf("encode data file: %w", err) @@ -108,6 +196,141 @@ func (s *Store) saveLocked() error { if err := os.Rename(temporaryName, s.path); err != nil { return fmt.Errorf("replace data file: %w", err) } + return syncDirectory(filepath.Dir(s.path)) +} + +func (s *Store) pruneLocked(now time.Time) { + if s.options.TaskRetention > 0 { + cutoff := now.Add(-s.options.TaskRetention) + for id, task := range s.state.Tasks { + if terminal(task.Status) && !task.UpdatedAt.IsZero() && task.UpdatedAt.Before(cutoff) { + delete(s.state.Tasks, id) + } + } + } + if s.options.EventRetention > 0 { + cutoff := now.Add(-s.options.EventRetention) + kept := s.state.Events[:0] + for _, event := range s.state.Events { + if event.ReceivedAt.IsZero() || !event.ReceivedAt.Before(cutoff) { + kept = append(kept, event) + } + } + s.state.Events = kept + } + if s.options.AuditRetention > 0 { + cutoff := now.Add(-s.options.AuditRetention) + kept := s.state.Audit[:0] + for _, entry := range s.state.Audit { + if entry.At.IsZero() || !entry.At.Before(cutoff) { + kept = append(kept, entry) + } + } + s.state.Audit = kept + } +} + +func (s *Store) backupLocked() error { + if s.options.BackupCount < 1 { + return nil + } + now := time.Now().UTC() + if s.options.BackupInterval > 0 && !s.lastBackup.IsZero() && now.Sub(s.lastBackup) < s.options.BackupInterval { + return nil + } + data, err := os.ReadFile(s.path) + if errors.Is(err, os.ErrNotExist) { + return nil + } + if err != nil { + return fmt.Errorf("read current data for backup: %w", err) + } + if err := os.MkdirAll(s.options.BackupDir, 0o700); err != nil { + return fmt.Errorf("create backup directory: %w", err) + } + name := fmt.Sprintf("%s.%d.json", filepath.Base(s.path), time.Now().UTC().UnixNano()) + backupPath := filepath.Join(s.options.BackupDir, name) + temporary, err := os.CreateTemp(s.options.BackupDir, ".wxagent-backup-*") + if err != nil { + return fmt.Errorf("create backup file: %w", err) + } + temporaryName := temporary.Name() + defer os.Remove(temporaryName) + if err := temporary.Chmod(0o600); err != nil { + _ = temporary.Close() + return fmt.Errorf("protect backup file: %w", err) + } + if _, err := temporary.Write(data); err != nil { + _ = temporary.Close() + return fmt.Errorf("write backup file: %w", err) + } + if err := temporary.Sync(); err != nil { + _ = temporary.Close() + return fmt.Errorf("sync backup file: %w", err) + } + if err := temporary.Close(); err != nil { + return fmt.Errorf("close backup file: %w", err) + } + if err := os.Rename(temporaryName, backupPath); err != nil { + return fmt.Errorf("publish backup file: %w", err) + } + if err := pruneBackups(s.options.BackupDir, filepath.Base(s.path), s.options.BackupCount); err != nil { + return fmt.Errorf("prune backups: %w", err) + } + if err := syncDirectory(s.options.BackupDir); err != nil { + return err + } + s.lastBackup = now + return nil +} + +func pruneBackups(directory, base string, keep int) error { + entries, err := os.ReadDir(directory) + if err != nil { + return err + } + type backup struct { + name string + when time.Time + } + backups := make([]backup, 0, len(entries)) + prefix := base + "." + for _, entry := range entries { + if entry.IsDir() || !strings.HasPrefix(entry.Name(), prefix) || !strings.HasSuffix(entry.Name(), ".json") { + continue + } + info, err := entry.Info() + if err != nil { + return err + } + backups = append(backups, backup{name: entry.Name(), when: info.ModTime()}) + } + sort.Slice(backups, func(i, j int) bool { + if backups[i].when.Equal(backups[j].when) { + return backups[i].name > backups[j].name + } + return backups[i].when.After(backups[j].when) + }) + if len(backups) <= keep { + return nil + } + for _, item := range backups[keep:] { + if err := os.Remove(filepath.Join(directory, item.name)); err != nil && !errors.Is(err, os.ErrNotExist) { + return err + } + } + return nil +} + +func syncDirectory(path string) error { + directory, err := os.Open(path) + if err != nil { + return fmt.Errorf("open directory for sync: %w", err) + } + defer directory.Close() + if err := directory.Sync(); err != nil && !errors.Is(err, syscall.EINVAL) { + return fmt.Errorf("sync directory: %w", err) + } return nil } diff --git a/control-plane/tls.go b/control-plane/tls.go new file mode 100644 index 0000000..f8ef835 --- /dev/null +++ b/control-plane/tls.go @@ -0,0 +1,94 @@ +package controlplane + +import ( + "crypto/sha256" + "crypto/tls" + "crypto/x509" + "encoding/hex" + "errors" + "fmt" + "net/http" + "os" + "strings" +) + +func newTLSConfig(config ServerConfig) (*tls.Config, error) { + if (config.TLSCertFile == "") != (config.TLSKeyFile == "") { + return nil, errors.New("TLS cert and key must be configured together") + } + if config.TLSCertFile == "" { + if config.MTLSClientCAFile != "" || config.MTLSRequireNodeCert || config.MTLSRevokedCertsFile != "" { + return nil, errors.New("mTLS settings require TLS cert and key") + } + return nil, nil + } + certificate, err := tls.LoadX509KeyPair(config.TLSCertFile, config.TLSKeyFile) + if err != nil { + return nil, fmt.Errorf("load TLS certificate: %w", err) + } + result := &tls.Config{ + MinVersion: tls.VersionTLS13, + Certificates: []tls.Certificate{certificate}, + } + if config.MTLSClientCAFile == "" { + if config.MTLSRequireNodeCert || config.MTLSRevokedCertsFile != "" { + return nil, errors.New("mTLS client CA is required for node certificates or revocation") + } + return result, nil + } + caBytes, err := os.ReadFile(config.MTLSClientCAFile) + if err != nil { + return nil, fmt.Errorf("read mTLS client CA: %w", err) + } + pool := x509.NewCertPool() + if !pool.AppendCertsFromPEM(caBytes) { + return nil, errors.New("mTLS client CA does not contain a PEM certificate") + } + result.ClientCAs = pool + // Web users may use HTTPS without a client certificate. Node routes enforce + // a certificate separately, so one TLS listener serves both audiences. + result.ClientAuth = tls.VerifyClientCertIfGiven + if config.MTLSRevokedCertsFile != "" { + if _, err := os.Stat(config.MTLSRevokedCertsFile); err != nil { + return nil, fmt.Errorf("check revoked certificate file: %w", err) + } + } + return result, nil +} + +func (s *Server) clientCertificateAllowed(r *http.Request) bool { + if r.TLS == nil || len(r.TLS.PeerCertificates) == 0 { + return !s.config.MTLSRequireNodeCert + } + if s.config.MTLSRevokedCertsFile == "" { + return true + } + revoked, err := readRevokedFingerprints(s.config.MTLSRevokedCertsFile) + if err != nil { + return false + } + fingerprint := sha256.Sum256(r.TLS.PeerCertificates[0].Raw) + _, isRevoked := revoked[hex.EncodeToString(fingerprint[:])] + return !isRevoked +} + +func readRevokedFingerprints(path string) (map[string]struct{}, error) { + data, err := os.ReadFile(path) + if err != nil { + return nil, err + } + result := make(map[string]struct{}) + for _, line := range strings.Split(string(data), "\n") { + line = strings.TrimSpace(strings.SplitN(line, "#", 2)[0]) + if line == "" { + continue + } + line = strings.ReplaceAll(strings.ToLower(line), ":", "") + decoded, err := hex.DecodeString(line) + if err != nil || len(decoded) != sha256.Size { + return nil, errors.New("revoked certificate list contains an invalid SHA-256 fingerprint") + } + result[line] = struct{}{} + } + return result, nil +} diff --git a/docs/WxAgent-远程协议草案-v1.0.md b/docs/WxAgent-远程协议草案-v1.0.md index 37054d5..3acc796 100644 --- a/docs/WxAgent-远程协议草案-v1.0.md +++ b/docs/WxAgent-远程协议草案-v1.0.md @@ -1,12 +1,12 @@ # WxAgent 远程协议草案 v1.0 > 对应:[`WxAgent-远程多节点控制与白名单数据上报开发计划.md`](./WxAgent-远程多节点控制与白名单数据上报开发计划.md) -> 状态:本地闭环已冻结;生产 TLS、节点 Token 签发和数据保留策略需在部署环境配置。 +> 状态:本地闭环已冻结;生产 TLS/mTLS 证书、外部 Secret、撤销文件、备份和保留策略需在部署环境配置。 ## 1. 传输和认证 - 协议版本:`v1`。 -- 节点只主动向控制面发起 HTTP 请求;生产地址必须使用 HTTPS。HTTP 仅允许 `127.0.0.1`/`::1` 本地集成验证。 +- 节点只主动向控制面发起 HTTP 请求;生产地址必须使用 HTTPS,节点可要求 mTLS 客户端证书。HTTP 仅允许显式配置的私有网络验证或 `127.0.0.1`/`::1` 本地集成验证。 - 节点请求使用 `Authorization: Bearer `。Token 按节点签发、可撤销,不写入日志和控制面持久化文件。 - Web 用户先调用 `POST /v1/auth/login`,使用返回的短期 Bearer 会话访问管理 API。未认证的业务请求返回 `401`。 - 每个请求可携带 `X-Correlation-Id`;服务端响应同名响应头和 JSON `correlation_id`。 diff --git a/docs/WxAgent-远程多节点控制与白名单数据上报开发计划.md b/docs/WxAgent-远程多节点控制与白名单数据上报开发计划.md index 67bf1ee..67814a8 100644 --- a/docs/WxAgent-远程多节点控制与白名单数据上报开发计划.md +++ b/docs/WxAgent-远程多节点控制与白名单数据上报开发计划.md @@ -348,7 +348,7 @@ UIA 事件/本地读取/任务结果/诊断 Agent 只有在同时配置并成功验证以下两项后,才允许连接中心、消费远程任务或上报白名单数据: - `authAddress`:中心认证/接入地址,必须使用 HTTPS 或等价安全传输。 -- `token`:节点专用认证 Token,按节点单独签发,可撤销、可轮换。 +- `token`:节点专用认证 Token,按节点单独签发;本项目不实现 Token 轮换,泄露时通过外部配置禁用对应节点并重启控制面。 示例: @@ -559,6 +559,6 @@ Agent 只有在同时配置并成功验证以下两项后,才允许连接中 - 已实现多节点配置入口:控制面支持 `WXAGENT_NODE_TOKENS` JSON map;单节点环境变量仍兼容。节点 `service.json` 可通过 `remoteConfigurationFile` 使用 CLI 维护的远程配置。 - 已实现消息事件的稳定会话标识要求:监听器只能在唯一 AutomationId 可确认时生成远程事件;本地事件总线不缓存消息正文,非白名单正文不入队、不持久化、不发送。 - 已补充远程只读查询任务:`POST /v1/reads/sessions`、`POST /v1/reads/contacts`、`POST /v1/reads/messages`;结果沿用任务租约/状态/幂等模型,节点读取前按本地已验证白名单过滤,消息查询仅返回当前 UI 可见历史。 -- 自动化证据:`npm run build --prefix control-plane/web`;`cd control-plane && go test ./... && go vet ./...`;Core 测试 166/166、Service 测试 20/20;`dotnet build WxAgent.sln -c Release -p:EnableWindowsTargeting=true`;`scripts/remote-control-smoke.sh` 均通过。 +- 自动化证据:`npm run build --prefix control-plane/web`;`cd control-plane && go test ./... && go vet ./...`;Core 测试 168/168、Service 测试 20/20;`dotnet build WxAgent.sln -c Release -p:EnableWindowsTargeting=true`;`scripts/remote-control-smoke.sh` 均通过。 - 真机 `windows-direct-20260911` 已在已登录、未锁定会话中注册并保持 `Online`;`send-text` 任务已完成 `Pending → Running → Succeeded`,并通过文件传输助手本地历史读回唯一标记确认。远程读取已验证会话 1 条、私聊联系人 1 条、群 1 条、消息 3 条,四个读取任务均为 `Succeeded`。 -- 当前不能标记为生产 R5 完成:控制面仍是单进程 JSON 文件/HTTP MVP,生产 HTTPS/mTLS、Token/证书撤销与轮换、外部密钥管理、规模压测、备份恢复和完整保留策略仍待完成。`/api/v1/sessions/open` 的导航能力仍可能返回 `CapabilityDisabled`;远程消息历史不是数据库全量历史。详细部署步骤见 [`WxAgent-远程控制面架构与部署说明-v1.0.md`](./WxAgent-远程控制面架构与部署说明-v1.0.md)。 +- 当前不能标记为生产 R5 完成:控制面已支持直接 TLS、节点 mTLS、客户端证书撤销文件、外部 Secret 文件、自动备份/保留、active/passive 数据锁和就绪检查;生产仍待完成外部 Secret/证书接入演练、异地备份恢复、规模压测和故障演练。本项目不实现 Token 轮换;泄露时由外部配置禁用对应节点并重启控制面。`/api/v1/sessions/open` 的导航能力仍可能返回 `CapabilityDisabled`;远程消息历史不是数据库全量历史。详细部署步骤见 [`WxAgent-远程控制面架构与部署说明-v1.0.md`](./WxAgent-远程控制面架构与部署说明-v1.0.md)。 diff --git a/docs/WxAgent-远程控制面架构与部署说明-v1.0.md b/docs/WxAgent-远程控制面架构与部署说明-v1.0.md index 9d7c918..33d2875 100644 --- a/docs/WxAgent-远程控制面架构与部署说明-v1.0.md +++ b/docs/WxAgent-远程控制面架构与部署说明-v1.0.md @@ -17,7 +17,10 @@ - 远程读取可见会话、联系人/群和当前 UI 可见消息历史; - 节点本地白名单过滤、稳定身份校验、结果范围化保存和断线补传; - React 管理页面登录、节点/任务/事件/审计查询; -- Gitea Actions 构建并推送控制面容器镜像。 +- Gitea Actions 构建并推送控制面容器镜像; +- 直接 TLS、节点 mTLS、客户端证书指纹撤销文件和外部 `*_FILE` Secret 注入; +- 数据文件 active/passive 单写锁、自动备份、保留期限和 `/readyz` 就绪检查; +- 有界的就绪/只读任务并发基线脚本。 当前明确不提供: @@ -25,7 +28,7 @@ - 控制面直接修改节点白名单; - 非白名单会话、完整微信数据库或数据库密钥上传; - 远程全量数据库历史读取;消息读取只来自当前微信 UI 可见历史; -- 生产级 HTTPS/mTLS 终止、在线 Token 撤销/轮换、HA 和规模压测; +- active-active 多实例、跨节点复制和托管式 Secret/PKI 服务; - `/api/v1/sessions/open` 的远程导航能力。当前该接口可能返回 `409 CapabilityDisabled`,不影响只读任务接口。 ## 2. 总体架构 @@ -158,6 +161,7 @@ Pending → Accepted → Running → Succeeded | 心跳超时 | 45 秒 | Web 查询节点时标记 `Offline` | | Web 会话 | 8 小时 | 控制面内存会话 | | 任务轮询窗口 | 0–30 秒 | `wait_seconds` | +| 备份间隔 | 5 分钟 | `WXAGENT_CONTROL_PLANE_BACKUP_INTERVAL` | | 事件正文 | 16 KiB | 节点协议上限 | | 任务载荷 | 64 KiB | 节点/中心共同限制 | | 任务结果 | 512 KiB | 含读取内容 | @@ -299,9 +303,11 @@ POST /v1/nodes/{node_id}/events ```text /data/control-plane-data.json +/data/control-plane-data.json.lock +/data/backups/control-plane-data.json..json ``` -`Store` 在单进程 Mutex 下执行读写,写入临时文件、`Sync` 后原子替换,文件权限为 `0600`。持久化对象包括: +`Store` 在进程内 Mutex 和数据文件旁的非阻塞 flock 下执行读写;同一数据卷同时只允许一个活动控制面实例。写入临时文件、`Sync` 后原子替换,并在替换父目录后同步目录,文件权限为 `0600`。按 `BackupInterval`(默认 5 分钟)把旧版本写入备份目录并按 `BackupCount` 删除旧备份;任务、事件和审计按配置的保留时长在下一次写入前清理,运行中任务不会被清理。持久化对象包括: ```text PersistedState @@ -311,7 +317,7 @@ PersistedState └── audit 登录、任务和节点操作审计 ``` -当前没有数据库事务、跨实例锁、自动 TTL、异地备份或 HA。升级/重建容器前必须备份 `/data/control-plane-data.json`,恢复时保持文件权限并确保只有一个控制面实例挂载该数据文件。 +当前没有数据库事务或 active-active HA;flock 提供的是共享卷上的 active/passive 单写保护。自动备份和按时长保留已实现,但异地复制、恢复演练和备份监控仍是部署门禁。升级/重建容器前应使用 `scripts/control-plane-data.sh backup`,恢复前停止活动实例并使用 `restore`,恢复时保持文件权限并确保只有一个控制面实例挂载该数据文件。 ### 6.2 Windows 节点 @@ -446,15 +452,17 @@ docker run -d \ curl --fail http://10.1.1.104:18090/healthz ``` -更推荐使用 Compose、systemd 或编排平台的 Secret 注入功能,避免 Token 出现在 `docker inspect` 的长期配置、进程列表和 shell 历史中。当前镜像支持环境变量认证,但尚未内置外部 Secret Provider。 +更推荐使用 Compose、systemd 或编排平台的 Secret 注入功能,避免 Token 出现在 `docker inspect` 的长期配置、进程列表和 shell 历史中。当前镜像支持 `*_FILE` Secret 文件注入;Secret 文件本身仍由部署平台负责保护和挂载。 -单节点兼容环境变量也可用: +单节点兼容环境变量也可用;对应的 `*_FILE` 形式优先从文件读取: ```text WXAGENT_NODE_ID -WXAGENT_NODE_TOKEN +WXAGENT_NODE_TOKEN / WXAGENT_NODE_TOKEN_FILE WXAGENT_WEB_USER -WXAGENT_WEB_PASSWORD +WXAGENT_WEB_PASSWORD / WXAGENT_WEB_PASSWORD_FILE +WXAGENT_NODE_TOKENS / WXAGENT_NODE_TOKENS_FILE +WXAGENT_WEB_USERS / WXAGENT_WEB_USERS_FILE ``` 多节点/多 Web 用户使用 JSON map: @@ -464,17 +472,18 @@ WXAGENT_NODE_TOKENS={"node-a":"token-a","node-b":"token-b"} WXAGENT_WEB_USERS={"admin":"password-a","auditor":"password-b"} ``` -### 8.3 生产 HTTPS/mTLS 门禁 +### 8.3 生产 HTTPS/mTLS -当前 Go 进程直接提供 HTTP。生产不得把该 HTTP 端口直接暴露到公网;正式部署至少应: +当前 Go 进程支持直接 TLS。生产不得把明文端口直接暴露到公网;正式部署至少应: -1. 在受控入口终止 HTTPS,并配置可信证书、强 TLS 和安全 Header; -2. 节点到入口使用 HTTPS; -3. 按节点管理、轮换和撤销 Token,或在边界与节点接入层启用 mTLS; -4. Token/密码放入外部密钥管理,不进入容器环境快照和日志; -5. 仅允许管理端、节点网段访问对应路径; -6. 对 `/data` 做加密备份、恢复演练和保留策略; -7. 完成失效证书、重放、中心重启、节点断线和并发规模测试。 +1. 配置 TLS 证书和私钥,服务端最低 TLS 1.3; +2. 配置签发节点证书的客户端 CA,并启用 `WXAGENT_MTLS_REQUIRE_NODE_CERT=true`;Web 管理用户仍可只使用 HTTPS; +3. 节点使用 PEM 客户端证书和私钥,服务端按 CA 验证; +4. 将已泄露或失效的客户端证书 DER-SHA256 指纹逐行写入撤销文件。可用 `scripts/revoke-client-certificate.sh` 原子追加;控制面每次节点认证重新读取该文件,撤销文件缺失/不可读时启动或认证失败。Docker Secret 更新后需重建控制面容器; +5. Token 和密码通过 `WXAGENT_NODE_TOKENS_FILE`、`WXAGENT_WEB_USERS_FILE` 等外部 Secret 文件注入,不进入镜像和控制面数据文件;本项目不实现 Token 轮换; +6. 仅允许管理端、节点网段访问对应路径; +7. 对 `/data` 做自动备份、异地复制、恢复演练和保留策略; +8. 完成失效证书、泄露响应、重放、中心重启、节点断线和并发规模测试。 `allowInsecureHttp` 只用于显式允许的私有 IP 内网验证;公网 HTTP 永远拒绝。它不是生产加密替代方案。 @@ -651,6 +660,7 @@ WxAgent.Host.exe smoke ```bash curl --fail http://10.1.1.104:18090/healthz +curl --fail http://10.1.1.104:18090/readyz curl --fail -c /tmp/wxagent.cookies \ -H 'Content-Type: application/json' \ @@ -714,6 +724,27 @@ curl --fail -H "$AUTH" -H 'Content-Type: application/json' \ 发送任务必须结合任务终态和微信本地历史读回确认;`Succeeded` 代表 Agent 已确认完成调用,不代替业务侧读回核验。 +### 10.4 并发基线 + +在 staging 控制面和已注册测试节点上执行;默认只创建 `read-sessions` 读取任务,不执行写任务: + +```bash +WXAGENT_SCALE_WEB_PASSWORD='' \ + scripts/control-plane-scale-smoke.py \ + --base-url http://127.0.0.1:8090 \ + --node-id node-a --account-id account-a \ + --requests 100 --workers 10 +``` + +只测就绪接口时不需要 Web 密码或节点参数: + +```bash +scripts/control-plane-scale-smoke.py --mode ready \ + --base-url https://control.example.com --requests 500 --workers 50 +``` + +脚本输出成功数、HTTP 状态分布、p50/p95;每次压测应保存版本、配置、机器规格、数据文件大小和结果,不把结果直接当作 1000 路生产验收。 + ## 11. 故障排查 | 现象 | 优先检查 | @@ -741,7 +772,7 @@ curl --fail -H "$AUTH" -H 'Content-Type: application/json' \ - 控制面 Docker 镜像构建成功,`/healthz` 返回 `ok/v1`; - Go `go test ./...`、`go vet ./...` 通过; -- Core 测试 `166/166` 通过;Service 测试 `20/20` 通过;完整 solution Release build 通过; +- Core 测试 `168/168` 通过;Service 测试 `20/20` 通过;完整 solution Release build 通过; - `scripts/remote-control-smoke.sh` 通过,覆盖注册、心跳、任务租约、幂等、结果回传、读取结果保留和拒绝路径; - Windows 节点 `windows-direct-20260911` 在已登录、未锁定会话中上线,注册 `read-sessions`、`read-contacts`、`read-messages` 能力; - 真机 `send-text` 任务 `Pending → Running → Succeeded`,并用文件传输助手本地历史读回唯一标记; @@ -749,9 +780,9 @@ curl --fail -H "$AUTH" -H 'Content-Type: application/json' \ 提交或部署新版本前必须重新执行相关测试。以下条件满足前不能标记生产完成: -- 控制面 HTTPS、节点 HTTPS/mTLS 和密钥外置; -- Token/证书轮换、撤销和泄露响应演练; -- 单进程 JSON 存储的备份恢复、并发写入和容量上限验证; +- 在实际部署环境启用控制面 HTTPS、节点 HTTPS/mTLS 和 Secret 文件挂载; +- 客户端证书撤销、泄露响应和重新签发演练(不包含 Token 轮换); +- 自动备份/保留、异地复制、恢复演练、active/passive 故障切换、并发写入和容量上限验证; - 节点/控制面重启、断网、租约过期和 `ResultUnconfirmed` 人工核对; - 真实目标规模、长时间运行、数据保留和日志泄漏扫描; - 当前微信版本变化后的 `doctor`、脱敏 UI 树和 smoke 回归。 @@ -759,7 +790,8 @@ curl --fail -H "$AUTH" -H 'Content-Type: application/json' \ ## 13. 相关文件 - `control-plane/Dockerfile`:控制面多阶段镜像构建; -- `control-plane/server.go`、`protocol.go`、`store.go`:HTTP 路由、协议和持久化; +- `control-plane/server.go`、`protocol.go`、`store.go`、`tls.go`:HTTP/TLS 路由、协议、持久化和证书撤销; +- `control-plane/compose.production.yml`:Docker Secret、TLS/mTLS 和保留配置的生产示例; - `control-plane/web/`:React 管理端; - `.gitea/workflows/build-web-image.yml`:Gitea 镜像发布; - `node-agent/WxAgent.Service/RemoteAgentHostedService.cs`:节点注册、轮询、任务执行和补传; @@ -767,4 +799,7 @@ curl --fail -H "$AUTH" -H 'Content-Type: application/json' \ - `node-agent/WxAgent.Core/RemoteTaskLedger.cs`:本地任务结果账本; - `node-agent/WxAgent.Host/RemoteCliCommands.cs`:远程认证、状态和白名单 CLI; - `docs/validation/Agent-Install-2026-09-11.md`:Windows 安装与真实验证记录; -- `scripts/remote-control-smoke.sh`:Linux/临时控制面闭环检查。 +- `scripts/remote-control-smoke.sh`:Linux/临时控制面闭环检查; +- `scripts/control-plane-scale-smoke.py`:就绪/只读任务并发基线; +- `scripts/control-plane-data.sh`:控制面数据备份和锁保护的原子恢复; +- `scripts/revoke-client-certificate.sh`:计算客户端证书指纹并原子追加撤销列表。 diff --git a/node-agent/WxAgent.Core/RemoteContracts.cs b/node-agent/WxAgent.Core/RemoteContracts.cs index 612eb57..178d63c 100644 --- a/node-agent/WxAgent.Core/RemoteContracts.cs +++ b/node-agent/WxAgent.Core/RemoteContracts.cs @@ -64,6 +64,18 @@ public sealed record RemoteAgentOptions [JsonPropertyName("token")] public string? Token { get; init; } + [JsonPropertyName("tokenFile")] + public string? TokenFile { get; init; } + + [JsonPropertyName("serverCaFile")] + public string? ServerCaFile { get; init; } + + [JsonPropertyName("clientCertificateFile")] + public string? ClientCertificateFile { get; init; } + + [JsonPropertyName("clientCertificateKeyFile")] + public string? ClientCertificateKeyFile { get; init; } + [JsonPropertyName("nodeId")] public string? NodeId { get; init; } @@ -74,15 +86,38 @@ public sealed record RemoteAgentOptions public bool AllowInsecureHttp { get; init; } public bool IsConfigured => !string.IsNullOrWhiteSpace(AuthAddress) - && !string.IsNullOrWhiteSpace(Token) + && (!string.IsNullOrWhiteSpace(Token) || !string.IsNullOrWhiteSpace(TokenFile)) && !string.IsNullOrWhiteSpace(NodeId); - public string TokenState => string.IsNullOrWhiteSpace(Token) ? "not-configured" : "configured"; + public string TokenState => !string.IsNullOrWhiteSpace(TokenFile) + ? "file-configured" + : string.IsNullOrWhiteSpace(Token) ? "not-configured" : "configured"; + + public string GetToken() + { + string? token; + if (string.IsNullOrWhiteSpace(TokenFile)) + { + token = Token; + } + else + { + try { token = File.ReadAllText(TokenFile).Trim(); } + catch (Exception exception) when (exception is IOException or UnauthorizedAccessException or ArgumentException) + { + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "tokenFile could not be read.", exception); + } + } + if (string.IsNullOrWhiteSpace(token)) + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "token or tokenFile must contain a non-empty token."); + return token; + } public void Validate() { - if (string.IsNullOrWhiteSpace(AuthAddress) || string.IsNullOrWhiteSpace(Token) || string.IsNullOrWhiteSpace(NodeId)) - throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "authAddress, token and nodeId are required before remote access is enabled."); + if (string.IsNullOrWhiteSpace(AuthAddress) || (!string.IsNullOrWhiteSpace(Token) && !string.IsNullOrWhiteSpace(TokenFile)) + || (!IsConfigured) || string.IsNullOrWhiteSpace(NodeId)) + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "authAddress, token or tokenFile and nodeId are required before remote access is enabled."); if (!Uri.TryCreate(AuthAddress, UriKind.Absolute, out var uri) || uri is null || uri.AbsolutePath == "/" && uri.Query.Length != 0 || uri.UserInfo.Length != 0 @@ -92,12 +127,28 @@ public sealed record RemoteAgentOptions && (!AllowInsecureHttp || !IsPrivateNetwork(uri.Host))) throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "Non-loopback HTTP requires allowInsecureHttp=true and a private-network IP address."); ValidateIdentifier(NodeId, "nodeId", 200); - if (Token.Any(char.IsWhiteSpace)) + var token = GetToken(); + if (token.Any(char.IsWhiteSpace)) throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "token must not contain whitespace."); + if (!string.IsNullOrWhiteSpace(TokenFile) && !File.Exists(TokenFile)) + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "tokenFile does not exist."); + ValidateOptionalFile(ServerCaFile, "serverCaFile"); + ValidateOptionalFile(ClientCertificateFile, "clientCertificateFile"); + ValidateOptionalFile(ClientCertificateKeyFile, "clientCertificateKeyFile"); + if (!string.IsNullOrWhiteSpace(ClientCertificateKeyFile) && string.IsNullOrWhiteSpace(ClientCertificateFile)) + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "clientCertificateFile is required with clientCertificateKeyFile."); + if (!string.IsNullOrWhiteSpace(ClientCertificateFile) && string.IsNullOrWhiteSpace(ClientCertificateKeyFile)) + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "clientCertificateKeyFile is required for PEM client certificates."); } public RemoteAgentOptions Redacted() => this with { Token = string.IsNullOrWhiteSpace(Token) ? null : "" }; + private static void ValidateOptionalFile(string? path, string name) + { + if (!string.IsNullOrWhiteSpace(path) && !File.Exists(path)) + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, $"{name} does not exist."); + } + private static bool IsLoopback(string host) => host.Equals("localhost", StringComparison.OrdinalIgnoreCase) || System.Net.IPAddress.TryParse(host.Trim('[', ']'), out var address) && System.Net.IPAddress.IsLoopback(address); diff --git a/node-agent/WxAgent.Core/RemoteControlClient.cs b/node-agent/WxAgent.Core/RemoteControlClient.cs index d170ae9..c58e6f3 100644 --- a/node-agent/WxAgent.Core/RemoteControlClient.cs +++ b/node-agent/WxAgent.Core/RemoteControlClient.cs @@ -1,4 +1,6 @@ using System.Net.Http.Headers; +using System.Net.Security; +using System.Security.Cryptography.X509Certificates; using System.Text; using System.Text.Json; @@ -18,13 +20,24 @@ public sealed class RemoteControlClient : IDisposable private readonly bool _ownsHttp; private readonly RemoteAgentOptions _options; private readonly Uri? _baseAddress; + private readonly X509Certificate2? _clientCertificate; + private readonly X509Certificate2? _serverCaCertificate; private bool _authenticated; public RemoteControlClient(RemoteAgentOptions options, HttpClient? httpClient = null) { _options = options ?? throw new ArgumentNullException(nameof(options)); _ownsHttp = httpClient is null; - _http = httpClient ?? new HttpClient(); + if (httpClient is null) + { + _clientCertificate = LoadClientCertificate(options); + _serverCaCertificate = LoadServerCaCertificate(options); + _http = CreateHttpClient(_clientCertificate, _serverCaCertificate); + } + else + { + _http = httpClient; + } if (Uri.TryCreate(options.AuthAddress, UriKind.Absolute, out var address)) _baseAddress = new Uri($"{address.Scheme}://{address.Authority}/", UriKind.Absolute); AuthState = options.IsConfigured ? RemoteAuthState.NotConfigured : RemoteAuthState.NotConfigured; @@ -174,7 +187,7 @@ public sealed class RemoteControlClient : IDisposable using var request = new HttpRequestMessage(method, new Uri(_baseAddress, path.TrimStart('/'))); request.Headers.TryAddWithoutValidation("X-Correlation-Id", NewCorrelationId()); if (authenticated || method == HttpMethod.Post && path == "/v1/nodes/register") - request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", _options.Token); + request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", _options.GetToken()); if (body is not null) { var bytes = JsonSerializer.SerializeToUtf8Bytes(body, RemoteJson.Options); @@ -220,7 +233,7 @@ public sealed class RemoteControlClient : IDisposable if (!_options.IsConfigured) { AuthState = RemoteAuthState.NotConfigured; - throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "Remote authAddress, token and nodeId must all be configured."); + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "Remote authAddress, token or tokenFile and nodeId must all be configured."); } _options.Validate(); } @@ -249,11 +262,42 @@ public sealed class RemoteControlClient : IDisposable return new RemoteApiError("RemoteRequestFailed", "Remote request failed.", correlationId); } + private static X509Certificate2? LoadClientCertificate(RemoteAgentOptions options) + { + if (string.IsNullOrWhiteSpace(options.ClientCertificateFile)) return null; + return X509Certificate2.CreateFromPemFile(options.ClientCertificateFile, options.ClientCertificateKeyFile!); + } + + private static X509Certificate2? LoadServerCaCertificate(RemoteAgentOptions options) => + string.IsNullOrWhiteSpace(options.ServerCaFile) ? null : new X509Certificate2(options.ServerCaFile); + + private static HttpClient CreateHttpClient(X509Certificate2? clientCertificate, X509Certificate2? serverCaCertificate) + { + var handler = new HttpClientHandler(); + if (clientCertificate is not null) handler.ClientCertificates.Add(clientCertificate); + if (serverCaCertificate is not null) + { + handler.ServerCertificateCustomValidationCallback = (_, certificate, _, errors) => + { + if (certificate is null || (errors & (SslPolicyErrors.RemoteCertificateNameMismatch | SslPolicyErrors.RemoteCertificateNotAvailable)) != 0) + return false; + using var chain = new X509Chain(); + chain.ChainPolicy.TrustMode = X509ChainTrustMode.CustomRootTrust; + chain.ChainPolicy.CustomTrustStore.Add(serverCaCertificate); + chain.ChainPolicy.RevocationMode = X509RevocationMode.NoCheck; + return chain.Build(new X509Certificate2(certificate)); + }; + } + return new HttpClient(handler) { Timeout = TimeSpan.FromSeconds(15) }; + } + private static string Escape(string value) => Uri.EscapeDataString(value); private static string NewCorrelationId() => Guid.NewGuid().ToString("N"); public void Dispose() { if (_ownsHttp) _http.Dispose(); + _clientCertificate?.Dispose(); + _serverCaCertificate?.Dispose(); } } diff --git a/node-agent/WxAgent.Host/Program.cs b/node-agent/WxAgent.Host/Program.cs index 458ba51..1287a4c 100644 --- a/node-agent/WxAgent.Host/Program.cs +++ b/node-agent/WxAgent.Host/Program.cs @@ -1033,6 +1033,8 @@ static void PrintHelp() => Console.WriteLine(""" WxAgent.Host commands: serve --config remote auth show|set|clear --config + auth set options: --address --token | --token-file --node + [--server-ca-file ] [--client-certificate-file --client-certificate-key-file ] remote status show --config [--data-dir ] remote probe run --config [--timeout 60] remote reporting show|enable|disable|account-add|account-enable|account-disable|allow|deny --config diff --git a/node-agent/WxAgent.Host/RemoteCliCommands.cs b/node-agent/WxAgent.Host/RemoteCliCommands.cs index 84cd9c2..235e96d 100644 --- a/node-agent/WxAgent.Host/RemoteCliCommands.cs +++ b/node-agent/WxAgent.Host/RemoteCliCommands.cs @@ -49,7 +49,11 @@ internal static class RemoteCliCommands var remote = new RemoteAgentOptions { AuthAddress = Required(args, "--address"), - Token = Required(args, "--token"), + Token = Option(args, "--token"), + TokenFile = Option(args, "--token-file"), + ServerCaFile = Option(args, "--server-ca-file"), + ClientCertificateFile = Option(args, "--client-certificate-file"), + ClientCertificateKeyFile = Option(args, "--client-certificate-key-file"), NodeId = Required(args, "--node"), ActiveAccountId = Option(args, "--active-account"), AllowInsecureHttp = Has(args, "--allow-insecure-http") @@ -71,7 +75,7 @@ internal static class RemoteCliCommands { var configuration = await RemoteNodeConfigurationStore.LoadAsync(path, cancellationToken); configuration.Remote.Validate(); - using var client = new RemoteControlClient(configuration.Remote, new HttpClient { Timeout = TimeSpan.FromSeconds(15) }); + using var client = new RemoteControlClient(configuration.Remote); var accounts = configuration.Reporting.Accounts.Select(account => new RemoteAccountSummary( account.AccountId, string.Equals(account.AccountId, configuration.Remote.ActiveAccountId, StringComparison.Ordinal), diff --git a/node-agent/WxAgent.Service/RemoteAgentHostedService.cs b/node-agent/WxAgent.Service/RemoteAgentHostedService.cs index 1c0614e..90ffa6c 100644 --- a/node-agent/WxAgent.Service/RemoteAgentHostedService.cs +++ b/node-agent/WxAgent.Service/RemoteAgentHostedService.cs @@ -38,8 +38,7 @@ public sealed class RemoteAgentHostedService( } var remote = configuredRemote; - using var http = new HttpClient { Timeout = TimeSpan.FromSeconds(15) }; - using var client = new RemoteControlClient(remote, http); + using var client = new RemoteControlClient(remote); var ledger = new RemoteTaskLedger(Path.Combine(options.DataDirectory, "remote-task-ledger.json")); var eventQueue = remoteQueue ?? new RemoteEventQueue(Path.Combine(options.DataDirectory, "remote-event-queue.json")); var retry = RetryDelay; diff --git a/scripts/control-plane-data.sh b/scripts/control-plane-data.sh new file mode 100755 index 0000000..7fc3157 --- /dev/null +++ b/scripts/control-plane-data.sh @@ -0,0 +1,42 @@ +#!/usr/bin/env bash +set -euo pipefail + +usage() { + echo "usage: $0 backup | restore " >&2 + exit 2 +} + +[ "$#" -eq 3 ] || usage +command=$1 +source=$2 +target=$3 + +case "$command" in + backup) + [ -f "$source" ] || { echo "data file does not exist: $source" >&2; exit 1; } + mkdir -p "$target" + name="$(basename "$source").$(date -u +%Y%m%dT%H%M%SZ).manual.json" + temporary="$target/.wxagent-backup-$$.tmp" + trap 'rm -f "$temporary"' EXIT + install -m 600 "$source" "$temporary" + mv -f "$temporary" "$target/$name" + sha256sum "$target/$name" + ;; + restore) + [ -f "$source" ] || { echo "backup file does not exist: $source" >&2; exit 1; } + directory=$(dirname "$target") + mkdir -p "$directory" + lock="$target.lock" + exec 9>"$lock" + flock -n 9 || { echo "data file is in use; stop the active control plane first" >&2; exit 1; } + temporary="$directory/.wxagent-restore-$$.tmp" + trap 'rm -f "$temporary"' EXIT + install -m 600 "$source" "$temporary" + python3 -m json.tool "$temporary" >/dev/null + mv -f "$temporary" "$target" + sha256sum "$target" + ;; + *) + usage + ;; +esac diff --git a/scripts/control-plane-scale-smoke.py b/scripts/control-plane-scale-smoke.py new file mode 100755 index 0000000..00e2a7b --- /dev/null +++ b/scripts/control-plane-scale-smoke.py @@ -0,0 +1,124 @@ +#!/usr/bin/env python3 +"""Bounded control-plane readiness/read-task concurrency smoke test.""" + +from __future__ import annotations + +import argparse +import concurrent.futures +import http.client +import json +import os +import ssl +import statistics +import sys +import time +import uuid +from urllib.parse import urlsplit + + +def request(base_url: str, path: str, method: str, body: object | None, headers: dict[str, str], timeout: float): + url = base_url.rstrip("/") + path + parts = urlsplit(url) + if parts.scheme.lower() not in {"http", "https"} or not parts.hostname or parts.username or parts.password: + raise ValueError("base URL must be an http/https URL without user information") + target = parts.path or "/" + if parts.query: + target += "?" + parts.query + data = None if body is None else json.dumps(body, ensure_ascii=False).encode() + connection: http.client.HTTPConnection + if parts.scheme.lower() == "https": + connection = http.client.HTTPSConnection(parts.hostname, parts.port, timeout=timeout, context=ssl.create_default_context()) + else: + connection = http.client.HTTPConnection(parts.hostname, parts.port, timeout=timeout) + started = time.perf_counter() + try: + connection.request(method, target, body=data, headers={"Accept": "application/json", **headers}) + response = connection.getresponse() + return response.status, time.perf_counter() - started, response.read() + except OSError as error: + return 599, time.perf_counter() - started, str(error).encode() + finally: + connection.close() + + +def percentile(values: list[float], p: float) -> float: + if not values: + return 0.0 + ordered = sorted(values) + index = min(len(ordered) - 1, round((len(ordered) - 1) * p)) + return ordered[index] + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--base-url", default=os.getenv("WXAGENT_CONTROL_PLANE_URL", "http://127.0.0.1:8090")) + parser.add_argument("--mode", choices=("ready", "reads"), default="reads") + parser.add_argument("--web-user", default=os.getenv("WXAGENT_SCALE_WEB_USER", "admin")) + parser.add_argument("--node-id", default=os.getenv("WXAGENT_SCALE_NODE_ID", "")) + parser.add_argument("--account-id", default=os.getenv("WXAGENT_SCALE_ACCOUNT_ID", "")) + parser.add_argument("--requests", type=int, default=100) + parser.add_argument("--workers", type=int, default=10) + parser.add_argument("--timeout", type=float, default=15.0) + args = parser.parse_args() + if args.requests < 1 or args.requests > 5000 or args.workers < 1 or args.workers > 100: + parser.error("requests must be 1..5000 and workers must be 1..100") + + headers: dict[str, str] = {} + if args.mode == "reads": + password = os.getenv("WXAGENT_SCALE_WEB_PASSWORD") + if not password or not args.node_id or not args.account_id: + parser.error("reads mode requires WXAGENT_SCALE_WEB_PASSWORD, --node-id and --account-id") + status, _, payload = request( + args.base_url, + "/v1/auth/login", + "POST", + {"username": args.web_user, "password": password}, + {"Content-Type": "application/json"}, + args.timeout, + ) + if status != 200: + print(json.dumps({"ok": False, "phase": "login", "status": status, "body": payload[:256].decode(errors="replace")})) + return 1 + try: + login = json.loads(payload) + access_token = login["access_token"] + except (UnicodeDecodeError, json.JSONDecodeError, KeyError, TypeError) as error: + print(json.dumps({"ok": False, "phase": "login-response", "error": str(error)})) + return 1 + headers = {"Authorization": "Bearer " + access_token, "Content-Type": "application/json"} + + def one(index: int) -> tuple[int, float]: + if args.mode == "ready": + status, elapsed, _ = request(args.base_url, "/readyz", "GET", None, {}, args.timeout) + else: + body = { + "node_id": args.node_id, + "account_id": args.account_id, + "idempotency_key": f"scale-read-{uuid.uuid4().hex}-{index}", + "limit": 1, + "offset": 0, + } + status, elapsed, _ = request(args.base_url, "/v1/reads/sessions", "POST", body, headers, args.timeout) + return status, elapsed + + with concurrent.futures.ThreadPoolExecutor(max_workers=args.workers) as executor: + results = list(executor.map(one, range(args.requests))) + elapsed = [item[1] for item in results] + success = sum(200 <= status < 300 for status, _ in results) + summary = { + "ok": success == args.requests, + "mode": args.mode, + "requests": args.requests, + "workers": args.workers, + "success": success, + "failed": args.requests - success, + "p50_ms": round(statistics.median(elapsed) * 1000, 2) if elapsed else 0, + "p95_ms": round(percentile(elapsed, 0.95) * 1000, 2), + "statuses": {str(status): sum(item[0] == status for item in results) for status, _ in results}, + } + print(json.dumps(summary, ensure_ascii=False, sort_keys=True)) + return 0 if summary["ok"] else 1 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/scripts/remote-control-smoke.sh b/scripts/remote-control-smoke.sh index 2e4399a..d274811 100755 --- a/scripts/remote-control-smoke.sh +++ b/scripts/remote-control-smoke.sh @@ -30,7 +30,7 @@ trap cleanup EXIT ) >"$log_file" 2>&1 & control_pid=$! -for _ in $(seq 1 60); do +for _ in $(seq 1 600); do if curl -fsS "http://127.0.0.1:${port}/healthz" >/dev/null 2>&1; then break; fi sleep 0.1 done diff --git a/scripts/revoke-client-certificate.sh b/scripts/revoke-client-certificate.sh new file mode 100755 index 0000000..74826ff --- /dev/null +++ b/scripts/revoke-client-certificate.sh @@ -0,0 +1,30 @@ +#!/usr/bin/env bash +set -euo pipefail + +if [ "$#" -ne 2 ]; then + echo "usage: $0 " >&2 + exit 2 +fi +certificate=$1 +list=$2 +[ -f "$certificate" ] || { echo "certificate does not exist: $certificate" >&2; exit 1; } + +fingerprint=$(openssl x509 -in "$certificate" -outform DER | sha256sum | awk '{print $1}') +[ "${#fingerprint}" -eq 64 ] || { echo "could not calculate certificate fingerprint" >&2; exit 1; } +mkdir -p "$(dirname "$list")" +lock="$list.lock" +exec 9>"$lock" +flock -x 9 +if [ -f "$list" ] && grep -Fqx "$fingerprint" "$list"; then + printf '%s\n' "$fingerprint" + exit 0 +fi +temporary="${list}.tmp.$$" +trap 'rm -f "$temporary"' EXIT +if [ -f "$list" ]; then + cat "$list" > "$temporary" +fi +printf '%s\n' "$fingerprint" >> "$temporary" +chmod 600 "$temporary" +mv -f "$temporary" "$list" +printf '%s\n' "$fingerprint" diff --git a/tests/node-agent/WxAgent.Core.Tests/RemoteControlTests.cs b/tests/node-agent/WxAgent.Core.Tests/RemoteControlTests.cs index 4155d7b..b31ab70 100644 --- a/tests/node-agent/WxAgent.Core.Tests/RemoteControlTests.cs +++ b/tests/node-agent/WxAgent.Core.Tests/RemoteControlTests.cs @@ -144,6 +144,52 @@ public sealed class RemoteAgentOptionsTests Assert.Throws(() => options.Validate()); } + + [Fact] + public void TokenFileIsAcceptedWithoutPuttingTheTokenInConfiguration() + { + var path = Path.Combine(Path.GetTempPath(), "wxagent-token-" + Guid.NewGuid().ToString("N")); + try + { + File.WriteAllText(path, "file-token\n"); + var options = new RemoteAgentOptions + { + AuthAddress = "https://control.example", + TokenFile = path, + NodeId = "node-1" + }; + + options.Validate(); + Assert.Equal("file-token", options.GetToken()); + Assert.Null(options.Redacted().Token); + Assert.Equal("file-configured", options.TokenState); + } + finally + { + if (File.Exists(path)) File.Delete(path); + } + } + + [Fact] + public void ClientCertificateConfigurationRequiresBothPemFiles() + { + var cert = Path.GetTempFileName(); + try + { + var options = new RemoteAgentOptions + { + AuthAddress = "https://control.example", + Token = "node-token", + NodeId = "node-1", + ClientCertificateFile = cert + }; + Assert.Throws(() => options.Validate()); + } + finally + { + File.Delete(cert); + } + } } public sealed class RemoteAccountContextTests