Bind isolated Dispatcher Agent and OSS deployment inputs

This commit is contained in:
2026-09-30 07:40:37 +08:00
parent e37a7371be
commit 5ed5409994
9 changed files with 308 additions and 56 deletions
@@ -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 加载。
+16 -1
View File
@@ -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
+16
View File
@@ -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
}
@@ -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
}
@@ -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)
}
})
}
}
+89
View File
@@ -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
}
+62
View File
@@ -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)
}
}
-55
View File
@@ -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
}
+54
View File
@@ -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
}