refactor(config): remove obsolete upstream Dispatcher file parser
This commit is contained in:
@@ -127,6 +127,7 @@
|
||||
- 旧通话流程入口:旧 `Execute`/`ExecuteWithCapture` 使用可缺省的旧 AI 快照,现已无 Agent 调用者;删除入口及专属测试,保留同一 `executeFlow` 媒体顺序实现。先在当前 `ExecuteApproved` 的测试中补齐首次媒体等待上限、三轮对话、开场发送失败不伪造播放、拒联后不发送回答与 ASR-only 实际采集区间,再迁移共用测试夹具;现行获批通话测试和全仓测试通过。未声称真实媒体链路通过。
|
||||
- 旧 AI 配置与供应商链路:旧 `Snapshot`/独立 MQ 配置响应/独立授权、旧可变 Mock 与旧供应商流水线已无现行执行者;删除专属实现和测试。当前获批执行仍使用原值 AI 模式及每通话不可变参数;将仍被当前路径使用的 `Mode` 与 `TurnResult` 移至 `internal/ai/turn.go`。先在当前 SDK 隔离测试补齐 LLM 失败与空回答候选不自动重试,再删除旧 SDK 专属测试;ASR-only、完整 AI、关键词、时限和媒体的当前测试继续通过。旧供应商测试的历史通过不能代替真实供应商联调。
|
||||
- 旧项目内合同装载:测试先复现 `contracts.Files` 将 v0.1–v0.4 的历史本地 Schema 仍嵌入现行二进制;现在仅嵌入获批当前 local Schema/示例和原样上游 v1 来源,移除可从代码访问旧 local 的 `ReadLocal`、旧命令/事件/配置验证层及专属测试,将并发 Schema 缓存保留于当前验证器。项目历史 Schema 文件与外部上游来源仍原样留在仓库,没有改写来源哈希;历史文件不再构成运行时回退。当前合同正反例、上游来源读取和全仓 Go 测试通过,未代表外部 SaaS 已签收。
|
||||
- 旧 Dispatcher 本地配置解析:旧 `LoadDispatcherFile` 独占读取上游 v1 `dispatcher-config.schema.json`,没有现行 Dispatcher CLI 调用者;删除实现与旧文件专属测试。现行入口仍从部署环境独立读取 Dispatcher 身份、SaaS 只读配置端点及受控 OSS 配置,并保留当前校验和错误报告。历史上游 Schema 未改,删除旧入口不代表真实配置来源已联调。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -1,135 +0,0 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/url"
|
||||
"os"
|
||||
"path"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/contracts"
|
||||
"github.com/santhosh-tekuri/jsonschema/v6"
|
||||
)
|
||||
|
||||
type dispatcherFile struct {
|
||||
SchemaVersion string `json:"schema_version"`
|
||||
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"`
|
||||
} `json:"oss"`
|
||||
}
|
||||
|
||||
var dispatcherFileSchema = sync.OnceValues(func() (*jsonschema.Schema, error) {
|
||||
data, err := contracts.Files.ReadFile("upstream/" + contracts.SourceCommit + "/dispatcher-config.schema.json")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
doc, err := jsonschema.UnmarshalJSON(bytes.NewReader(data))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
compiler := jsonschema.NewCompiler()
|
||||
compiler.AssertFormat()
|
||||
const resource = "https://go-sip.local/dispatcher-config.json"
|
||||
if err := compiler.AddResource(resource, doc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return compiler.Compile(resource)
|
||||
})
|
||||
|
||||
// LoadDispatcherFile applies an entirely validated configuration atomically.
|
||||
// Only credentials explicitly referenced by this file are read from the
|
||||
// environment; missing/invalid files never fall back to legacy OSS settings.
|
||||
func (c *Config) LoadDispatcherFile(filename string) error {
|
||||
if c == nil || filename == "" {
|
||||
return errors.New("dispatcher configuration and --config file are required")
|
||||
}
|
||||
file, err := os.Open(filename)
|
||||
if err != nil {
|
||||
return fmt.Errorf("open dispatcher configuration: %w", err)
|
||||
}
|
||||
const maxBytes = 64 << 10
|
||||
raw, readErr := io.ReadAll(io.LimitReader(file, maxBytes+1))
|
||||
if err := errors.Join(readErr, file.Close()); err != nil {
|
||||
return fmt.Errorf("read dispatcher configuration: %w", err)
|
||||
}
|
||||
if len(raw) > maxBytes || !utf8.Valid(raw) {
|
||||
return errors.New("dispatcher configuration must be valid UTF-8 JSON of at most 64 KiB")
|
||||
}
|
||||
check := json.NewDecoder(bytes.NewReader(raw))
|
||||
if err := uniqueJSONKeys(check, 0); err != nil {
|
||||
return fmt.Errorf("dispatcher configuration JSON: %w", err)
|
||||
}
|
||||
if _, err := check.Token(); !errors.Is(err, io.EOF) {
|
||||
return errors.New("dispatcher configuration must contain one JSON object")
|
||||
}
|
||||
var input dispatcherFile
|
||||
decoder := json.NewDecoder(bytes.NewReader(raw))
|
||||
decoder.DisallowUnknownFields()
|
||||
if err := decoder.Decode(&input); err != nil {
|
||||
return fmt.Errorf("decode dispatcher configuration: %w", err)
|
||||
}
|
||||
schema, err := dispatcherFileSchema()
|
||||
if err != nil {
|
||||
return fmt.Errorf("compile dispatcher configuration schema: %w", err)
|
||||
}
|
||||
document, err := jsonschema.UnmarshalJSON(bytes.NewReader(raw))
|
||||
if err != nil {
|
||||
return fmt.Errorf("parse dispatcher configuration: %w", err)
|
||||
}
|
||||
if err := schema.Validate(document); err != nil {
|
||||
var failure *jsonschema.ValidationError
|
||||
if errors.As(err, &failure) {
|
||||
for len(failure.Causes) > 0 {
|
||||
failure = failure.Causes[0]
|
||||
}
|
||||
// Report the violated field and rule, never its potentially sensitive value.
|
||||
return fmt.Errorf("dispatcher configuration schema violation at /%s (%T)", strings.Join(failure.InstanceLocation, "/"), failure.ErrorKind)
|
||||
}
|
||||
return errors.New("dispatcher configuration schema validation failed")
|
||||
}
|
||||
endpoint, err := url.Parse(input.OSS.Endpoint)
|
||||
if err != nil || endpoint.Hostname() == "" || (endpoint.Scheme != "https" && endpoint.Scheme != "http") || endpoint.User != nil || endpoint.RawQuery != "" || endpoint.Fragment != "" {
|
||||
return errors.New("oss.endpoint must be an HTTP(S) service URL without credentials, query or fragment")
|
||||
}
|
||||
for name, value := range map[string]string{"region": input.OSS.Region, "bucket": input.OSS.Bucket, "object_prefix": input.OSS.ObjectPrefix} {
|
||||
if value == "" || strings.TrimSpace(value) != value {
|
||||
return fmt.Errorf("oss.%s must be nonempty and have no surrounding whitespace", name)
|
||||
}
|
||||
}
|
||||
prefix := input.OSS.ObjectPrefix
|
||||
if path.IsAbs(prefix) || path.Clean(prefix) != prefix || prefix == "." || prefix == ".." || strings.HasPrefix(prefix, "../") || strings.Contains(prefix, `\`) {
|
||||
return errors.New("oss.object_prefix must be a relative object prefix without traversal")
|
||||
}
|
||||
key, err := fileCredential("access_key_id_env", input.OSS.AccessKeyIDEnv)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
secret, err := fileCredential("access_key_secret_env", input.OSS.AccessKeySecretEnv)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
updated := *c
|
||||
updated.DispatcherID = input.DispatcherID
|
||||
updated.OSSEndpoint = input.OSS.Endpoint
|
||||
updated.OSSRegion = input.OSS.Region
|
||||
updated.OSSBucket = input.OSS.Bucket
|
||||
updated.OSSKeyPrefix = prefix
|
||||
updated.OSSAccessKeyID = key
|
||||
updated.OSSAccessKeySecret = secret
|
||||
updated.OSSGrantTTL = 15 * time.Minute
|
||||
*c = updated
|
||||
return nil
|
||||
}
|
||||
@@ -1,139 +0,0 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func dispatcherFileFixture(t *testing.T) map[string]any {
|
||||
t.Helper()
|
||||
t.Setenv("MQ_TEST_OSS_KEY", "test-key-not-a-real-credential")
|
||||
t.Setenv("MQ_TEST_OSS_SECRET", "test-secret-not-a-real-credential")
|
||||
return map[string]any{
|
||||
"schema_version": "1.0",
|
||||
"dispatcher_id": "c046b893-8628-4589-ae50-619d049248a6",
|
||||
"oss": map[string]any{
|
||||
"endpoint": "https://oss.example.invalid", "region": "test-region",
|
||||
"bucket": "test-bucket", "object_prefix": "recordings",
|
||||
"access_key_id_env": "MQ_TEST_OSS_KEY", "access_key_secret_env": "MQ_TEST_OSS_SECRET",
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func writeDispatcherFile(t *testing.T, raw []byte) string {
|
||||
t.Helper()
|
||||
path := filepath.Join(t.TempDir(), "dispatcher.json")
|
||||
if err := os.WriteFile(path, raw, 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return path
|
||||
}
|
||||
|
||||
func marshalDispatcherFile(t *testing.T, value map[string]any) []byte {
|
||||
t.Helper()
|
||||
raw, err := json.Marshal(value)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return raw
|
||||
}
|
||||
|
||||
func TestLoadDispatcherFile(t *testing.T) {
|
||||
value := dispatcherFileFixture(t)
|
||||
t.Setenv("DISPATCHER_OSS_ENDPOINT", "https://ignored.example.invalid")
|
||||
cfg := Config{Mode: "mock", DBPath: "preserved.db", OSSGrantTTL: time.Hour, OSSBucket: "old"}
|
||||
if err := cfg.LoadDispatcherFile(writeDispatcherFile(t, marshalDispatcherFile(t, value))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if cfg.DispatcherID != value["dispatcher_id"] || cfg.OSSEndpoint != "https://oss.example.invalid" || cfg.OSSRegion != "test-region" || cfg.OSSBucket != "test-bucket" || cfg.OSSKeyPrefix != "recordings" || cfg.OSSGrantTTL != 15*time.Minute {
|
||||
t.Fatalf("wrong file mapping: id=%q endpoint=%q region=%q bucket=%q prefix=%q ttl=%v", cfg.DispatcherID, cfg.OSSEndpoint, cfg.OSSRegion, cfg.OSSBucket, cfg.OSSKeyPrefix, cfg.OSSGrantTTL)
|
||||
}
|
||||
if cfg.OSSAccessKeyID != os.Getenv("MQ_TEST_OSS_KEY") || cfg.OSSAccessKeySecret != os.Getenv("MQ_TEST_OSS_SECRET") {
|
||||
t.Fatal("explicit credential references not resolved")
|
||||
}
|
||||
if cfg.Mode != "mock" || cfg.DBPath != "preserved.db" {
|
||||
t.Fatal("unrelated deployment settings overwritten")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadDispatcherFileRejectsInvalidValuesAtomically(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
mutate func(map[string]any)
|
||||
}{
|
||||
{"version", func(m map[string]any) { m["schema_version"] = "2.0" }},
|
||||
{"missing ID", func(m map[string]any) { delete(m, "dispatcher_id") }},
|
||||
{"invalid ID", func(m map[string]any) { m["dispatcher_id"] = "dispatcher" }},
|
||||
{"unknown root", func(m map[string]any) { m["legacy"] = true }},
|
||||
{"case alias", func(m map[string]any) { m["SCHEMA_VERSION"] = m["schema_version"]; delete(m, "schema_version") }},
|
||||
{"nested case alias", func(m map[string]any) {
|
||||
m["oss"].(map[string]any)["BUCKET"] = "test-bucket"
|
||||
delete(m["oss"].(map[string]any), "bucket")
|
||||
}},
|
||||
{"missing oss", func(m map[string]any) { delete(m, "oss") }},
|
||||
{"plaintext secret", func(m map[string]any) { m["oss"].(map[string]any)["access_key_secret"] = "not-a-real-secret" }},
|
||||
{"TTL override", func(m map[string]any) { m["oss"].(map[string]any)["grant_ttl_seconds"] = 3600 }},
|
||||
{"missing bucket", func(m map[string]any) { delete(m["oss"].(map[string]any), "bucket") }},
|
||||
{"blank region", func(m map[string]any) { m["oss"].(map[string]any)["region"] = " " }},
|
||||
{"invalid endpoint", func(m map[string]any) { m["oss"].(map[string]any)["endpoint"] = "file:///tmp/oss" }},
|
||||
{"URL credentials", func(m map[string]any) {
|
||||
m["oss"].(map[string]any)["endpoint"] = "https://user:password@example.invalid"
|
||||
}},
|
||||
{"empty prefix", func(m map[string]any) { m["oss"].(map[string]any)["object_prefix"] = "" }},
|
||||
{"prefix escape", func(m map[string]any) { m["oss"].(map[string]any)["object_prefix"] = "recordings/../other" }},
|
||||
{"invalid env reference", func(m map[string]any) { m["oss"].(map[string]any)["access_key_id_env"] = "${SECRET}" }},
|
||||
{"missing credential", func(m map[string]any) { m["oss"].(map[string]any)["access_key_secret_env"] = "MQ_TEST_MISSING" }},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
value := dispatcherFileFixture(t)
|
||||
t.Setenv("MQ_TEST_MISSING", "")
|
||||
tc.mutate(value)
|
||||
before := Config{DispatcherID: "unchanged", OSSAccessKeySecret: "untouched"}
|
||||
cfg := before
|
||||
err := cfg.LoadDispatcherFile(writeDispatcherFile(t, marshalDispatcherFile(t, value)))
|
||||
if err == nil {
|
||||
t.Fatal("invalid file accepted")
|
||||
}
|
||||
if cfg != before {
|
||||
t.Fatal("invalid file partially changed configuration")
|
||||
}
|
||||
if strings.Contains(err.Error(), os.Getenv("MQ_TEST_OSS_SECRET")) {
|
||||
t.Fatal("credential leaked in error")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadDispatcherFileRejectsDuplicateAndTrailingJSON(t *testing.T) {
|
||||
raw := string(marshalDispatcherFile(t, dispatcherFileFixture(t)))
|
||||
for _, bad := range []string{
|
||||
strings.Replace(raw, `"schema_version":"1.0"`, `"schema_version":"1.0","schema_version":"1.0"`, 1),
|
||||
strings.Replace(raw, `"bucket":"test-bucket"`, `"bucket":"test-bucket","bucket":"other"`, 1),
|
||||
raw + ` {}`, raw + ` null`, raw[:len(raw)-1], `null`, `[]`,
|
||||
strings.Repeat("[", 100) + strings.Repeat("]", 100),
|
||||
strings.Repeat(" ", 65537), string([]byte{0xff}),
|
||||
} {
|
||||
var cfg Config
|
||||
if err := cfg.LoadDispatcherFile(writeDispatcherFile(t, []byte(bad))); err == nil {
|
||||
t.Fatal("invalid/ambiguous JSON accepted")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadDispatcherFileDoesNotFallBack(t *testing.T) {
|
||||
var cfg Config
|
||||
t.Setenv("DISPATCHER_OSS_BUCKET", "legacy-bucket")
|
||||
for _, path := range []string{"", filepath.Join(t.TempDir(), "missing.json"), t.TempDir()} {
|
||||
if err := cfg.LoadDispatcherFile(path); err == nil {
|
||||
t.Fatal("missing/invalid file accepted")
|
||||
}
|
||||
}
|
||||
var absent *Config
|
||||
if err := absent.LoadDispatcherFile("unused"); err == nil {
|
||||
t.Fatal("nil config accepted")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user