diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index be8b946..1a4c245 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -34,6 +34,7 @@ - 新增无实现代次字段的 Dispatcher 环境预检:必须显式给出规范 UUID v4 身份、只读 HTTP 地址及密钥、RabbitMQ 地址和 SQLite 路径;当前 Mock 仅接受本机 HTTP/MQ 目标,mixed/real 在任何资源操作前拒绝。预检不打开数据库或网络,缺失值不继承旧默认配置,错误不打印凭据。`go test ./internal/config -run '^TestLoadDispatcherEnvironment' -count=1` 通过;此预检尚未接入主 CLI,不能当作 P02/P03 完成。 - Dispatcher 另有纯运行环境预检:唯一 Agent 端点清单和 OSS 配置文件须提供路径,本机 gRPC 监听及双向 TLS 文件/Agent 证书指纹须显式给出;缺失、非法指纹或非本机监听明确拒绝且不回显值。该检查不读取配置文件、不连接 Agent、不打开 SQLite/MQ;`go test ./internal/config -run '^TestLoadDispatcherRuntimeEnvironment' -count=1` 通过。当前仍未由 `dispatcher` 主命令调用,不能当作 Agent 加载或 OSS 签发已验收。 +- D 的部署文件读取已独立于旧代次 Schema:Agent 清单仅接受**一个本机端点**;OSS 配置必须匹配本 D 身份、指定本机 HTTPS 目标、15 分钟授权与显式资产上限,凭据只经受控环境变量引用。超大文件、重复/未知字段、多个 Agent、非本机目标、缺失凭据和错误 D 归属均拒绝;JSON 唯一键及凭据引用基础函数已移出旧配置文件。`go test ./internal/config -count=1` 通过。此处仅验证部署数据,不代表 Agent 实际加载、官方 SDK 已签发目标,亦未接入主 `dispatcher` 命令。 - Agent 增加独立的无代次环境预检:Agent/Cell 与预授权 D 的身份、会话与私有恢复路径、隔离 Mock 场景和已加载 SIP 测试事实、D gRPC 目标及 mTLS 文件/指纹均须显式提供;仅允许本机监听和本机 D 端点,mixed/real 在访问文件或网络前拒绝。缺失项不继承旧 `FromEnv` 默认值,凭据/地址不回显;`go test ./internal/config -run '^TestLoadAgentEnvironment' -count=1` 通过。此检查已接入根命令可达的唯一 Agent Mock 入口;本机双向 TLS 监听可由受信 D 证书探测,另一张同 CA 证书被指纹门禁拒绝。本地测试用批准的 D 身份激活会话,错误 D UUID 即使携带受信证书也在写入会话日志前拒绝;尚未执行呼叫、验证 Agent→D 录音或真实 Asterisk 加载。 diff --git a/internal/config/agent_endpoints.go b/internal/config/agent_endpoints.go index bdd7023..c6dff7c 100644 --- a/internal/config/agent_endpoints.go +++ b/internal/config/agent_endpoints.go @@ -26,10 +26,25 @@ func LoadAgentEndpoints(path string) ([]AgentEndpoint, error) { if strings.TrimSpace(path) == "" { return nil, nil } - data, err := os.ReadFile(path) + file, err := os.Open(path) if err != nil { return nil, fmt.Errorf("read Agent endpoint inventory: %w", err) } + const maxBytes = 64 << 10 + data, readErr := io.ReadAll(io.LimitReader(file, maxBytes+1)) + if err := errors.Join(readErr, file.Close()); err != nil { + return nil, fmt.Errorf("read Agent endpoint inventory: %w", err) + } + if len(data) > maxBytes { + return nil, errors.New("Agent endpoint inventory exceeds 64 KiB") + } + check := json.NewDecoder(bytes.NewReader(data)) + if err := uniqueJSONKeys(check, 0); err != nil { + return nil, fmt.Errorf("Agent endpoint inventory contains invalid or duplicate JSON: %w", err) + } + if _, err := check.Token(); !errors.Is(err, io.EOF) { + return nil, errors.New("Agent endpoint inventory contains trailing JSON") + } decoder := json.NewDecoder(bytes.NewReader(data)) decoder.DisallowUnknownFields() var endpoints []AgentEndpoint diff --git a/internal/config/credential_ref.go b/internal/config/credential_ref.go new file mode 100644 index 0000000..675c599 --- /dev/null +++ b/internal/config/credential_ref.go @@ -0,0 +1,16 @@ +package config + +import ( + "fmt" + "os" + "strings" +) + +// Deployment files identify credential environment variables, not secret bytes. +func fileCredential(field, reference string) (string, error) { + value, exists := os.LookupEnv(reference) + if !exists || strings.TrimSpace(value) == "" { + return "", fmt.Errorf("credential referenced by oss.%s is unavailable", field) + } + return value, nil +} diff --git a/internal/config/current_agent_endpoints.go b/internal/config/current_agent_endpoints.go new file mode 100644 index 0000000..7eb0a02 --- /dev/null +++ b/internal/config/current_agent_endpoints.go @@ -0,0 +1,25 @@ +package config + +import ( + "errors" + "strings" +) + +// LoadCurrentMockAgentEndpoint accepts exactly the one deployment-owned local +// Agent assigned to the isolated Dispatcher. SaaS task data cannot override it. +func LoadCurrentMockAgentEndpoint(filename string) (AgentEndpoint, error) { + if strings.TrimSpace(filename) == "" { + return AgentEndpoint{}, errors.New("DISPATCHER_AGENT_ENDPOINTS_FILE is required") + } + endpoints, err := LoadAgentEndpoints(filename) + if err != nil { + return AgentEndpoint{}, err + } + if len(endpoints) != 1 { + return AgentEndpoint{}, errors.New("current Mock Dispatcher requires exactly one assigned Agent endpoint") + } + if !localGRPCAddress(endpoints[0].Address, false) { + return AgentEndpoint{}, errors.New("current Mock Agent endpoint must be local with an explicit port") + } + return endpoints[0], nil +} diff --git a/internal/config/current_agent_endpoints_test.go b/internal/config/current_agent_endpoints_test.go new file mode 100644 index 0000000..4fbc495 --- /dev/null +++ b/internal/config/current_agent_endpoints_test.go @@ -0,0 +1,45 @@ +package config + +import ( + "os" + "path/filepath" + "strings" + "testing" +) + +const mockAgentEndpoint = `[{"agent_id":"agent-mock","cell_id":"cell-mock","address":"127.0.0.1:19090","server_name":"agent.local"}]` + +func writeMockAgentEndpoints(t *testing.T, body string) string { + t.Helper() + filename := filepath.Join(t.TempDir(), "agent-endpoints.json") + if err := os.WriteFile(filename, []byte(body), 0600); err != nil { + t.Fatal(err) + } + return filename +} + +func TestLoadCurrentMockAgentEndpointRequiresOneLocalDeploymentTarget(t *testing.T) { + endpoint, err := LoadCurrentMockAgentEndpoint(writeMockAgentEndpoints(t, mockAgentEndpoint)) + if err != nil || endpoint.AgentID != "agent-mock" || endpoint.CellID != "cell-mock" || + endpoint.Address != "127.0.0.1:19090" || endpoint.ServerName != "agent.local" { + t.Fatalf("approved local Agent target was changed: %v", err) + } + if _, err := LoadCurrentMockAgentEndpoint(""); err == nil { + t.Fatal("missing target inventory was defaulted") + } + for _, tc := range []struct { + name, body string + }{ + {"multiple Agents", `[{"agent_id":"agent-mock","cell_id":"cell-mock","address":"127.0.0.1:19090","server_name":"agent.local"},{"agent_id":"other","cell_id":"other","address":"127.0.0.1:19091","server_name":"other.local"}]`}, + {"external address", strings.Replace(mockAgentEndpoint, "127.0.0.1:19090", "saas.example.invalid:443", 1)}, + {"public listener", strings.Replace(mockAgentEndpoint, "127.0.0.1:19090", "0.0.0.0:19090", 1)}, + {"duplicate address", strings.Replace(mockAgentEndpoint, `"address":"127.0.0.1:19090"`, `"address":"127.0.0.1:19090","address":"127.0.0.1:19091"`, 1)}, + {"oversized inventory", strings.Repeat(" ", 65<<10) + mockAgentEndpoint}, + } { + t.Run(tc.name, func(t *testing.T) { + if _, err := LoadCurrentMockAgentEndpoint(writeMockAgentEndpoints(t, tc.body)); err == nil || strings.Contains(err.Error(), "saas.example.invalid") { + t.Fatalf("unapproved Agent target accepted or echoed: %v", err) + } + }) + } +} diff --git a/internal/config/current_oss.go b/internal/config/current_oss.go new file mode 100644 index 0000000..1fd2d25 --- /dev/null +++ b/internal/config/current_oss.go @@ -0,0 +1,89 @@ +package config + +import ( + "bytes" + "encoding/json" + "errors" + "io" + "net/url" + "os" + "path" + "strings" + "time" + "unicode/utf8" + + "git.ipao.vip/rogee/go-sip/internal/oss" + "git.ipao.vip/rogee/go-sip/internal/tenant" +) + +// LoadCurrentOSSConfig reads only the deployment-owned OSS settings for this +// Dispatcher. Credentials are resolved from explicit environment references; +// the file, its contents and secrets are never echoed in errors. +func LoadCurrentOSSConfig(filename, dispatcherID string) (oss.Config, error) { + if strings.TrimSpace(filename) == "" || tenant.ValidateDispatcherID(dispatcherID) != nil { + return oss.Config{}, errors.New("DISPATCHER_OSS_CONFIG_FILE and canonical Dispatcher ID are required") + } + file, err := os.Open(filename) + if err != nil { + return oss.Config{}, errors.New("DISPATCHER_OSS_CONFIG_FILE cannot be opened") + } + const maxBytes = 64 << 10 + raw, readErr := io.ReadAll(io.LimitReader(file, maxBytes+1)) + if err := errors.Join(readErr, file.Close()); err != nil || len(raw) > maxBytes || !utf8.Valid(raw) { + return oss.Config{}, errors.New("DISPATCHER_OSS_CONFIG_FILE must be readable UTF-8 JSON within 64 KiB") + } + check := json.NewDecoder(bytes.NewReader(raw)) + if err := uniqueJSONKeys(check, 0); err != nil { + return oss.Config{}, errors.New("DISPATCHER_OSS_CONFIG_FILE has invalid or repeated fields") + } + if _, err := check.Token(); !errors.Is(err, io.EOF) { + return oss.Config{}, errors.New("DISPATCHER_OSS_CONFIG_FILE must contain one JSON object") + } + var input struct { + DispatcherID string `json:"dispatcher_id"` + OSS struct { + Endpoint string `json:"endpoint"` + Region string `json:"region"` + Bucket string `json:"bucket"` + ObjectPrefix string `json:"object_prefix"` + AccessKeyIDEnv string `json:"access_key_id_env"` + AccessKeySecretEnv string `json:"access_key_secret_env"` + MaxAssetBytes int64 `json:"max_asset_bytes"` + } `json:"oss"` + } + decoder := json.NewDecoder(bytes.NewReader(raw)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&input); err != nil { + return oss.Config{}, errors.New("DISPATCHER_OSS_CONFIG_FILE has invalid fields") + } + if input.DispatcherID != dispatcherID { + return oss.Config{}, errors.New("OSS configuration is not assigned to this Dispatcher") + } + endpoint, err := url.Parse(input.OSS.Endpoint) + if err != nil || !localEndpoint(input.OSS.Endpoint, "https") || endpoint.User != nil || endpoint.RawQuery != "" || endpoint.Fragment != "" || + (endpoint.Path != "" && endpoint.Path != "/") || endpoint.RawPath != "" || !localGRPCAddress(endpoint.Host, false) { + return oss.Config{}, errors.New("OSS Mock endpoint must be a local HTTPS service with an explicit port") + } + prefix := input.OSS.ObjectPrefix + if prefix == "" || strings.TrimSpace(prefix) != prefix || path.IsAbs(prefix) || path.Clean(prefix) != prefix || + prefix == "." || prefix == ".." || strings.HasPrefix(prefix, "../") || strings.Contains(prefix, `\`) { + return oss.Config{}, errors.New("OSS object prefix must be a relative path without traversal") + } + key, err := fileCredential("access_key_id_env", input.OSS.AccessKeyIDEnv) + if err != nil { + return oss.Config{}, err + } + secret, err := fileCredential("access_key_secret_env", input.OSS.AccessKeySecretEnv) + if err != nil { + return oss.Config{}, err + } + configuration := oss.Config{ + Endpoint: input.OSS.Endpoint, Region: input.OSS.Region, Bucket: input.OSS.Bucket, + KeyPrefix: prefix, AccessKeyID: key, AccessKeySecret: secret, + GrantTTL: 15 * time.Minute, MaxAssetBytes: input.OSS.MaxAssetBytes, + } + if err := configuration.Validate(); err != nil { + return oss.Config{}, err + } + return configuration, nil +} diff --git a/internal/config/current_oss_test.go b/internal/config/current_oss_test.go new file mode 100644 index 0000000..49eb022 --- /dev/null +++ b/internal/config/current_oss_test.go @@ -0,0 +1,62 @@ +package config + +import ( + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +const mockOSSDeployment = `{"dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","oss":{"endpoint":"https://127.0.0.1:19445","region":"cn-mock","bucket":"mock-bucket","object_prefix":"approved","access_key_id_env":"MOCK_OSS_KEY_ID","access_key_secret_env":"MOCK_OSS_KEY_SECRET","max_asset_bytes":1048576}}` + +func writeMockOSSDeployment(t *testing.T, body string) string { + t.Helper() + path := filepath.Join(t.TempDir(), "oss.json") + if err := os.WriteFile(path, []byte(body), 0600); err != nil { + t.Fatal(err) + } + return path +} + +func TestLoadCurrentOSSConfigRequiresBoundLocalHTTPSAndExplicitCredentials(t *testing.T) { + t.Setenv("MOCK_OSS_KEY_ID", "isolated-id") + t.Setenv("MOCK_OSS_KEY_SECRET", "isolated-secret") + configuration, err := LoadCurrentOSSConfig(writeMockOSSDeployment(t, mockOSSDeployment), dispatcherFixtureID) + if err != nil { + t.Fatal(err) + } + if configuration.Endpoint != "https://127.0.0.1:19445" || configuration.Bucket != "mock-bucket" || + configuration.KeyPrefix != "approved" || configuration.MaxAssetBytes != 1048576 || configuration.GrantTTL != 15*time.Minute || + configuration.AccessKeyID != "isolated-id" || configuration.AccessKeySecret != "isolated-secret" { + t.Fatal("approved deployment settings were rewritten") + } +} + +func TestLoadCurrentOSSConfigRejectsUnapprovedDeploymentWithoutLeakingCredentials(t *testing.T) { + t.Setenv("MOCK_OSS_KEY_ID", "isolated-id") + t.Setenv("MOCK_OSS_KEY_SECRET", "isolated-secret") + for _, tc := range []struct { + name string + body string + }{ + {"wrong Dispatcher", strings.Replace(mockOSSDeployment, dispatcherFixtureID, "22222222-2222-4222-8222-222222222222", 1)}, + {"legacy generation", strings.Replace(mockOSSDeployment, `{"dispatcher_id":`, `{"schema_version":"v3","dispatcher_id":`, 1)}, + {"insecure endpoint", strings.Replace(mockOSSDeployment, "https://127.0.0.1:19445", "http://127.0.0.1:19445", 1)}, + {"external endpoint", strings.Replace(mockOSSDeployment, "https://127.0.0.1:19445", "https://oss.example.invalid:443", 1)}, + {"duplicate key", strings.Replace(mockOSSDeployment, `"bucket":"mock-bucket"`, `"bucket":"mock-bucket","bucket":"other"`, 1)}, + {"missing size", strings.Replace(mockOSSDeployment, `"max_asset_bytes":1048576`, `"max_asset_bytes":0`, 1)}, + {"extra document", mockOSSDeployment + ` {}`}, + {"oversized", strings.Repeat(" ", 65<<10) + mockOSSDeployment}, + } { + t.Run(tc.name, func(t *testing.T) { + if _, err := LoadCurrentOSSConfig(writeMockOSSDeployment(t, tc.body), dispatcherFixtureID); err == nil || strings.Contains(err.Error(), "isolated-secret") { + t.Fatalf("unapproved deployment admitted or exposed credential: %v", err) + } + }) + } + t.Setenv("MOCK_OSS_KEY_SECRET", "") + if _, err := LoadCurrentOSSConfig(writeMockOSSDeployment(t, mockOSSDeployment), dispatcherFixtureID); err == nil || strings.Contains(err.Error(), "isolated-id") { + t.Fatalf("unavailable credential was defaulted or exposed: %v", err) + } +} diff --git a/internal/config/dispatcher_file.go b/internal/config/dispatcher_file.go index 7dc6859..2543ae6 100644 --- a/internal/config/dispatcher_file.go +++ b/internal/config/dispatcher_file.go @@ -133,58 +133,3 @@ func (c *Config) LoadDispatcherFile(filename string) error { *c = updated return nil } - -func fileCredential(field, reference string) (string, error) { - value, exists := os.LookupEnv(reference) - if !exists || strings.TrimSpace(value) == "" { - return "", fmt.Errorf("credential referenced by oss.%s is unavailable", field) - } - return value, nil -} - -// encoding/json accepts repeated object keys. Inspect its token stream before -// typed decoding so duplicate settings cannot silently override earlier values. -func uniqueJSONKeys(decoder *json.Decoder, depth int) error { - if depth > 16 { - return errors.New("configuration nesting is too deep") - } - token, err := decoder.Token() - if err != nil { - return err - } - delimiter, compound := token.(json.Delim) - if !compound { - return nil - } - switch delimiter { - case '{': - seen := make(map[string]bool) - for decoder.More() { - keyToken, err := decoder.Token() - if err != nil { - return err - } - key, ok := keyToken.(string) - if !ok { - return errors.New("configuration object key must be a string") - } - if seen[key] { - return fmt.Errorf("duplicate configuration field %q", key) - } - seen[key] = true - if err := uniqueJSONKeys(decoder, depth+1); err != nil { - return err - } - } - case '[': - for decoder.More() { - if err := uniqueJSONKeys(decoder, depth+1); err != nil { - return err - } - } - default: - return errors.New("unexpected configuration JSON delimiter") - } - _, err = decoder.Token() - return err -} diff --git a/internal/config/json_keys.go b/internal/config/json_keys.go new file mode 100644 index 0000000..511b2de --- /dev/null +++ b/internal/config/json_keys.go @@ -0,0 +1,54 @@ +package config + +import ( + "encoding/json" + "errors" + "fmt" +) + +// encoding/json accepts repeated object keys. Inspect the token stream before +// typed decoding so deployment settings cannot silently override earlier values. +func uniqueJSONKeys(decoder *json.Decoder, depth int) error { + if depth > 16 { + return errors.New("configuration nesting is too deep") + } + token, err := decoder.Token() + if err != nil { + return err + } + delimiter, compound := token.(json.Delim) + if !compound { + return nil + } + switch delimiter { + case '{': + seen := make(map[string]bool) + for decoder.More() { + keyToken, err := decoder.Token() + if err != nil { + return err + } + key, ok := keyToken.(string) + if !ok { + return errors.New("configuration object key must be a string") + } + if seen[key] { + return fmt.Errorf("duplicate configuration field %q", key) + } + seen[key] = true + if err := uniqueJSONKeys(decoder, depth+1); err != nil { + return err + } + } + case '[': + for decoder.More() { + if err := uniqueJSONKeys(decoder, depth+1); err != nil { + return err + } + } + default: + return errors.New("unexpected configuration JSON delimiter") + } + _, err = decoder.Token() + return err +}