feat(test): serve contract-bound SaaS fixtures and provision isolated MQ
This commit is contained in:
@@ -81,6 +81,7 @@
|
||||
- 唯一现行 SaaS↔Dispatcher 业务规范是 [`docs/thirds/saas-dispatcher.md`](docs/thirds/saas-dispatcher.md);当前项目内 Schema、拓扑、正反例及来源/hash 在 [`contracts/local/`](contracts/local/);内部 Agent RPC 在 [`proto/agent/agent.proto`](proto/agent/agent.proto)。本地验收与外部缺口见 [`docs/evidence/saas-dispatcher-p08-acceptance.md`](docs/evidence/saas-dispatcher-p08-acceptance.md)。Markdown 不代替机器合同或外部签收,也不另外维护一份平行字段定义。
|
||||
- 已由当前合同来源清单固定哈希的历史提案与计划保留**原字节**于 [`docs/archive/sources/`](docs/archive/sources/README.md),使用者原有未提交的两份旧对接文档也按原字节归档;旧上游 v1 在 [`docs/archive/upstream/`](docs/archive/upstream/README.md) 可离线校验,但不嵌入运行合同。旧 F/W 工作包、旧 MQ-only 合同及归档不作为当前运行入口。固定 MQ `v1`、HTTP `/internal/v1/dispatcher/...` 和业务 revision 是现行通信规则,不是自有实现代次;不得为历史路径新建兼容或回退。
|
||||
- 当前业务范围仍仅**单节点、单 Dispatcher、单 Agent、单 Cell、单租户和隔离 Mock**。根命令只接受显式 `agent`/`dispatcher`;mixed/real 启动即拒绝。另有严格隔离的 `--mode sip-only`:只允许 Dispatcher 读完整 SIP、持久接纳归属 `sip.config`、经已激活双向 TLS Agent 会话将完整快照应用到原生 Asterisk,并从运行态核对版本;不发现任务、不启动业务呼叫、不开放准入、不处理其他业务控制。测试机已凭使用者批准,将三条历史登记地址以**测试快照显式声明的** UDP/IP 鉴权、无需 REGISTER 配置写入并核对 Asterisk 运行态;供应商尚未确认这些实际线路是否满足上述参数,无拨号,不证明线路可用或生产签收。没有真实 SaaS、management、真实通话或生产签收;真实百炼 LLM 与 OSS 已各做一次**不拨号、独立的最小连接/写入诊断**(见 [`docs/evidence/real-ai-oss-one-shot-20261003.md`](docs/evidence/real-ai-oss-one-shot-20261003.md)),这不是获批任务的 ASR/LLM/TTS、录音和上传全链路签收。生产发布包仍为 `production_approval=false`。任何本机 Mock 或 SIP-only 核验均不授权真实呼叫。
|
||||
- 使用者批准独立的测试专用 `deploys/test/saas-mock/` 仅模拟缺失的 SaaS:从 0600 的静态快照提供五类正式 HTTP 配置,在专用 RabbitMQ vhost 中由模拟 SaaS 预建现行拓扑;不在 Dispatcher 内注入快照,不发布外呼消息,不替代 Agent/Asterisk/AI/OSS,也不解除当前 Mock-only 呼叫屏障。只有完成真实链路、线路鉴权和拨号前诊断,才可另行安排逐次试拨。此模拟不构成真实 SaaS 签收。
|
||||
- 开发按 TDD 分批,小步提交;不得覆盖使用者现存修改/未跟踪文件,不自动清理、迁移或覆盖任何现存 SQLite、spool、outbox 和 Agent 恢复文件。旧 `.executions` 及恢复根目录中旧 `.uploads`、`.upload-locks`、逐执行 `state.json` 的发现须只读失败关闭,现存未交付事实由使用者确认处置。真实云账号、EIP、线路、拨号、生产部署和共享数据操作分别需要明确授权。
|
||||
|
||||
## SaaS、Dispatcher 与 Agent 的现行边界
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
# SaaS 侧测试模拟服务(非生产)
|
||||
|
||||
仅替代当前缺失的 SaaS,不改 Dispatcher 的正式 HTTP 配置接口或 MQ 归属。**当前只提供静态配置和预建队列;不会投递 `call.execute`、拨号或产生业务结果。** 虚构的合同示例不能充当真实 AI、线路或任务授权。
|
||||
|
||||
## 数据
|
||||
|
||||
准备仅自己可读的目录,包含 `sip.json`、`providers.json`、`quota.json` 和 `tasks/*.json`;所有 JSON 文件须为普通文件且权限为 `0600`。分别对应 [`contracts/local/`](../../../contracts/local/) 的 `sip_config`、`ai_providers`、`tenant_quota`、`task_config`;每份快照的 `dispatcher_id` 必须相同,任务须属于同一租户且文件名为 `<task_id>.json`。本阶段最多六项任务,启动时全部校验并读入内存;更改文件后须重新启动,不热替换在途任务。现有 `contracts/local/examples/` **仅用于隔离 Mock 测试**,不得直接复制成真实拨号授权。
|
||||
|
||||
## 运行
|
||||
|
||||
- 事先建立专用空 RabbitMQ vhost,名称以 `saas-mock-` 开头;模拟服务只创建现行 durable 交换机、控制队列、每任务队列与结果队列并做精确绑定,不自动清理/覆盖已有消息。Dispatcher 自身仍只被动核验拓扑。
|
||||
- 准备测试 HTTPS 证书与私钥;将 `SAAS_MOCK_DISPATCHER_SECRET` 和含凭据的 `SAAS_MOCK_RABBITMQ_URL` 放在受限环境文件,不在命令行、仓库或聊天中传输。
|
||||
- 启动:`go run ./deploys/test/saas-mock --data <私有目录> --dispatcher-id <UUID> --listen <地址:端口> --tls-cert <证书文件> --tls-key <私钥文件>`。
|
||||
- 服务以标准 `X-DISPATCHER-id` 和 `X-DISPATCHER-SECRET-KEY` 校验 Dispatcher,再提供五类只读配置。错误的归属、资源、租户、快照或消息队列会导致拒绝启动/读取。
|
||||
|
||||
测试:`go test ./deploys/test/saas-mock` 验证正式配置客户端;设置指向**单独隔离 vhost** 的 `SAAS_MOCK_TEST_BROKER_URL` 后,`TestSaaSMockProvisionsDispatcherTopology` 还将实际预建 MQ 并用 Dispatcher 被动读回。缺省测试不会连接共享 RabbitMQ。
|
||||
|
||||
当前正式呼叫仍只允许 Mock,测试机 ARI/HTTP、真实 AI/TTS、录音、OSS 与拨号前抓包尚未完成验证。六组外呼均未执行。此服务通过测试也不代表真实 SaaS 或生产签收。
|
||||
@@ -0,0 +1,36 @@
|
||||
// saas-mock serves approved, static SaaS test snapshots through the formal
|
||||
// Dispatcher configuration API. It never creates calls or publishes MQ jobs.
|
||||
package main
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"flag"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
func main() {
|
||||
dataDir := flag.String("data", "", "private directory with sip.json, providers.json, quota.json and tasks/*.json")
|
||||
dispatcherID := flag.String("dispatcher-id", "", "approved test Dispatcher UUID")
|
||||
listen := flag.String("listen", "", "explicit HTTPS listen address")
|
||||
cert := flag.String("tls-cert", "", "HTTPS certificate file")
|
||||
key := flag.String("tls-key", "", "HTTPS private key file")
|
||||
flag.Parse()
|
||||
secret := os.Getenv("SAAS_MOCK_DISPATCHER_SECRET")
|
||||
brokerURL := os.Getenv("SAAS_MOCK_RABBITMQ_URL")
|
||||
if *dataDir == "" || *dispatcherID == "" || *listen == "" || *cert == "" || *key == "" || secret == "" || brokerURL == "" || flag.NArg() != 0 {
|
||||
log.Fatal(errors.New("data, dispatcher-id, listen, tls-cert, tls-key, SAAS_MOCK_DISPATCHER_SECRET and SAAS_MOCK_RABBITMQ_URL are required"))
|
||||
}
|
||||
data, err := loadDataset(*dataDir, *dispatcherID)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
if err := provisionMQ(brokerURL, data); err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
log.Printf("SaaS test configuration and MQ topology ready: tasks=%d dispatcher=%s", len(data.tasks), *dispatcherID)
|
||||
server := http.Server{Addr: *listen, Handler: data.handler(secret), ReadHeaderTimeout: 5 * time.Second}
|
||||
log.Fatal(server.ListenAndServeTLS(*cert, *key))
|
||||
}
|
||||
@@ -0,0 +1,147 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
)
|
||||
|
||||
const testDispatcher = "c046b893-8628-4589-ae50-619d049248a6"
|
||||
|
||||
func testDataDir(t *testing.T) string {
|
||||
t.Helper()
|
||||
root := t.TempDir()
|
||||
if err := os.Mkdir(filepath.Join(root, "tasks"), 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for target, source := range map[string]string{
|
||||
"sip.json": "config-read-sip.json",
|
||||
"providers.json": "config-read-providers.json",
|
||||
"quota.json": "config-read-quota.json",
|
||||
"tasks/task-full.json": "config-read-task-full.json",
|
||||
} {
|
||||
original, err := os.ReadFile(filepath.Join("..", "..", "..", "contracts", "local", "examples", source))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(root, target), original, 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
return root
|
||||
}
|
||||
|
||||
func TestSaaSMockServesFormalReadContract(t *testing.T) {
|
||||
data, err := loadDataset(testDataDir(t), testDispatcher)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
server := httptest.NewServer(data.handler("test-only-secret"))
|
||||
defer server.Close()
|
||||
client, err := configread.NewClient(server.URL, testDispatcher, "test-only-secret", server.Client())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx := context.Background()
|
||||
sip, err := client.ReadSIP(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
providers, err := client.ReadProviders(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
tasks, cursor, err := client.ReadAllTasks(ctx)
|
||||
if err != nil || len(tasks) != 1 || tasks[0].TaskID != "task-full" || cursor == "" {
|
||||
t.Fatalf("formal task discovery: count=%d cursor=%q err=%v", len(tasks), cursor, err)
|
||||
}
|
||||
snapshot, err := client.ReadTask(ctx, "task-full", 1001, sip, providers)
|
||||
if err != nil || snapshot.Task.TaskID != "task-full" || snapshot.Quota.TenantID != snapshot.Task.TenantID {
|
||||
t.Fatalf("formal task/quota read: task=%q quota=%d err=%v", snapshot.Task.TaskID, snapshot.Quota.TenantID, err)
|
||||
}
|
||||
if _, err := client.ReadTask(ctx, "absent", 1001, sip, providers); err == nil {
|
||||
t.Fatal("unknown task must fail closed")
|
||||
}
|
||||
wrong, err := configread.NewClient(server.URL, testDispatcher, "wrong-secret", server.Client())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := wrong.ReadSIP(ctx); err == nil || !strings.Contains(err.Error(), "HTTP 403") {
|
||||
t.Fatalf("wrong dispatcher identity should be rejected: %v", err)
|
||||
}
|
||||
request, err := http.NewRequest(http.MethodPost, server.URL+"/internal/v1/dispatcher/sip", nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
request.Header.Set("X-DISPATCHER-id", testDispatcher)
|
||||
request.Header.Set("X-DISPATCHER-SECRET-KEY", "test-only-secret")
|
||||
response, err := server.Client().Do(request)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer response.Body.Close()
|
||||
if response.StatusCode != http.StatusMethodNotAllowed {
|
||||
t.Fatalf("configuration must stay read-only: HTTP %d", response.StatusCode)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSaaSMockDiscoversSixDistinctTasks(t *testing.T) {
|
||||
root := testDataDir(t)
|
||||
original, err := os.ReadFile(filepath.Join(root, "tasks", "task-full.json"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for n := 2; n <= 6; n++ {
|
||||
var task map[string]any
|
||||
if err := json.Unmarshal(original, &task); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
id := fmt.Sprintf("task-%d", n)
|
||||
task["task_id"] = id
|
||||
body, err := json.Marshal(task)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(root, "tasks", id+".json"), body, 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
data, err := loadDataset(root, testDispatcher)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
server := httptest.NewServer(data.handler("test-only-secret"))
|
||||
defer server.Close()
|
||||
client, err := configread.NewClient(server.URL, testDispatcher, "test-only-secret", server.Client())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
tasks, cursor, err := client.ReadAllTasks(context.Background())
|
||||
if err != nil || len(tasks) != 6 || cursor != "mock-complete" {
|
||||
t.Fatalf("six independent formal tasks were not discovered: count=%d cursor=%q err=%v", len(tasks), cursor, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSaaSMockRejectsInvalidOwnerAtStartup(t *testing.T) {
|
||||
root := testDataDir(t)
|
||||
path := filepath.Join(root, "tasks", "task-full.json")
|
||||
original, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
bad := strings.Replace(string(original), testDispatcher, "00000000-0000-4000-8000-000000000000", 1)
|
||||
if err := os.WriteFile(path, []byte(bad), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := loadDataset(root, testDispatcher); err == nil {
|
||||
t.Fatal("a task assigned to another dispatcher cannot be served")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"strings"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/mq"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// provisionMQ is SaaS-owned setup in a dedicated vhost. It never publishes,
|
||||
// consumes, purges, deletes, or resets an existing queue.
|
||||
func provisionMQ(brokerURL string, data dataset) error {
|
||||
address, err := url.Parse(brokerURL)
|
||||
if err != nil || (address.Scheme != "amqp" && address.Scheme != "amqps") || address.Host == "" ||
|
||||
!strings.HasPrefix(strings.TrimPrefix(address.Path, "/"), "saas-mock-") {
|
||||
return errors.New("SaaS test MQ requires a dedicated saas-mock-* vhost")
|
||||
}
|
||||
conn, err := amqp.Dial(brokerURL)
|
||||
if err != nil {
|
||||
return fmt.Errorf("connect SaaS test MQ: %T", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
return fmt.Errorf("open SaaS test MQ channel: %T", err)
|
||||
}
|
||||
defer channel.Close()
|
||||
for _, name := range []string{mq.CommandsExchange, mq.ResultsExchange, mq.DeadLetterExchange} {
|
||||
if err := channel.ExchangeDeclare(name, "topic", true, false, false, false, nil); err != nil {
|
||||
return fmt.Errorf("declare SaaS test exchange %s: %T", name, err)
|
||||
}
|
||||
}
|
||||
control := mq.ControlQueueName(data.dispatcherID)
|
||||
if _, err := channel.QueueDeclare(control, true, false, false, false, nil); err != nil {
|
||||
return fmt.Errorf("declare SaaS control queue: %T", err)
|
||||
}
|
||||
if err := channel.QueueBind(control, "d."+data.dispatcherID+".control.in", mq.CommandsExchange, false, nil); err != nil {
|
||||
return fmt.Errorf("bind SaaS control queue: %T", err)
|
||||
}
|
||||
for _, task := range data.discovery {
|
||||
queue := "agent-call.d." + data.dispatcherID + ".task." + task.TaskID + ".v1"
|
||||
if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
||||
return fmt.Errorf("declare SaaS task queue: %T", err)
|
||||
}
|
||||
if err := channel.QueueBind(queue, "d."+data.dispatcherID+".task."+task.TaskID+".in", mq.CommandsExchange, false, nil); err != nil {
|
||||
return fmt.Errorf("bind SaaS task queue: %T", err)
|
||||
}
|
||||
}
|
||||
const resultQueue = "agent-call.saas.events.v1"
|
||||
if _, err := channel.QueueDeclare(resultQueue, true, false, false, false, nil); err != nil {
|
||||
return fmt.Errorf("declare SaaS result queue: %T", err)
|
||||
}
|
||||
if err := channel.QueueBind(resultQueue, "d."+data.dispatcherID+".out", mq.ResultsExchange, false, nil); err != nil {
|
||||
return fmt.Errorf("bind SaaS result queue: %T", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/mq"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
func TestSaaSMockProvisionRejectsSharedVHost(t *testing.T) {
|
||||
data, err := loadDataset(testDataDir(t), testDispatcher)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := provisionMQ("amqp://guest:test-only@127.0.0.1:5672/", data); err == nil || !strings.Contains(err.Error(), "dedicated") {
|
||||
t.Fatalf("SaaS simulator must not touch the default/shared vhost: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSaaSMockProvisionsDispatcherTopology(t *testing.T) {
|
||||
brokerURL := os.Getenv("SAAS_MOCK_TEST_BROKER_URL")
|
||||
if brokerURL == "" {
|
||||
t.Skip("requires an explicitly provisioned, private saas-mock-* RabbitMQ vhost")
|
||||
}
|
||||
data, err := loadDataset(testDataDir(t), testDispatcher)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := provisionMQ(brokerURL, data); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
broker, err := mq.Open(brokerURL, testDispatcher, 1)
|
||||
if err != nil {
|
||||
t.Fatalf("Dispatcher must be able to verify the SaaS-owned topology: %v", err)
|
||||
}
|
||||
defer broker.Close()
|
||||
conn, err := amqp.Dial(brokerURL)
|
||||
if err != nil {
|
||||
t.Fatal("read back SaaS-provisioned task queue")
|
||||
}
|
||||
defer conn.Close()
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer channel.Close()
|
||||
if _, err := channel.QueueDeclarePassive("agent-call.d."+testDispatcher+".task.task-full.v1", true, false, false, false, nil); err != nil {
|
||||
t.Fatalf("task queue not precreated by SaaS simulator: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,162 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"crypto/subtle"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
"git.ipao.vip/rogee/go-sip/internal/contract"
|
||||
)
|
||||
|
||||
type dataset struct {
|
||||
dispatcherID string
|
||||
tenantID int64
|
||||
sip []byte
|
||||
providers []byte
|
||||
quota []byte
|
||||
tasks map[string][]byte
|
||||
discovery []configread.DiscoveredTask
|
||||
}
|
||||
|
||||
func loadDataset(dir, dispatcherID string) (dataset, error) {
|
||||
if dispatcherID == "" {
|
||||
return dataset{}, errors.New("mock dispatcher ID is required")
|
||||
}
|
||||
read := func(name, resource string) ([]byte, error) {
|
||||
path := filepath.Join(dir, name)
|
||||
info, err := os.Stat(path)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read SaaS test snapshot %s: %w", name, err)
|
||||
}
|
||||
if !info.Mode().IsRegular() || info.Mode().Perm()&0077 != 0 {
|
||||
return nil, fmt.Errorf("SaaS test snapshot %s must be a private regular file", name)
|
||||
}
|
||||
body, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read SaaS test snapshot %s: %w", name, err)
|
||||
}
|
||||
if err := contract.ValidateCurrent("config-read", body); err != nil {
|
||||
// Schema errors can contain credential values; never include their text.
|
||||
return nil, fmt.Errorf("SaaS test snapshot %s violates the current contract", name)
|
||||
}
|
||||
var owner struct {
|
||||
DispatcherID string `json:"dispatcher_id"`
|
||||
Resource string `json:"resource"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &owner); err != nil || owner.DispatcherID != dispatcherID || owner.Resource != resource {
|
||||
return nil, fmt.Errorf("SaaS test snapshot %s has a wrong owner or resource", name)
|
||||
}
|
||||
return body, nil
|
||||
}
|
||||
data := dataset{dispatcherID: dispatcherID, tasks: make(map[string][]byte)}
|
||||
var err error
|
||||
if data.sip, err = read("sip.json", "sip_config"); err != nil {
|
||||
return dataset{}, err
|
||||
}
|
||||
if data.providers, err = read("providers.json", "ai_providers"); err != nil {
|
||||
return dataset{}, err
|
||||
}
|
||||
if data.quota, err = read("quota.json", "tenant_quota"); err != nil {
|
||||
return dataset{}, err
|
||||
}
|
||||
var quota configread.Quota
|
||||
if err = json.Unmarshal(data.quota, "a); err != nil || quota.TenantID <= 0 {
|
||||
return dataset{}, errors.New("SaaS test quota has no tenant ID")
|
||||
}
|
||||
data.tenantID = quota.TenantID
|
||||
files, err := filepath.Glob(filepath.Join(dir, "tasks", "*.json"))
|
||||
if err != nil || len(files) == 0 || len(files) > 6 {
|
||||
return dataset{}, errors.New("SaaS test dataset must contain one to six task snapshots")
|
||||
}
|
||||
for _, path := range files {
|
||||
body, err := read(filepath.Join("tasks", filepath.Base(path)), "task_config")
|
||||
if err != nil {
|
||||
return dataset{}, err
|
||||
}
|
||||
var task configread.Task
|
||||
if err := json.Unmarshal(body, &task); err != nil || task.TaskID == "" || filepath.Base(path) != task.TaskID+".json" || task.TenantID != data.tenantID {
|
||||
return dataset{}, errors.New("SaaS test task has an invalid ID or tenant")
|
||||
}
|
||||
if _, found := data.tasks[task.TaskID]; found {
|
||||
return dataset{}, errors.New("duplicate SaaS test task ID")
|
||||
}
|
||||
data.tasks[task.TaskID] = body
|
||||
data.discovery = append(data.discovery, configread.DiscoveredTask{
|
||||
TaskID: task.TaskID, TenantID: task.TenantID, TaskRevision: task.TaskRevision, Status: task.Status,
|
||||
})
|
||||
}
|
||||
sort.Slice(data.discovery, func(i, j int) bool { return data.discovery[i].TaskID < data.discovery[j].TaskID })
|
||||
page, err := json.Marshal(configread.TaskPage{DispatcherID: dispatcherID, Cursor: "mock-complete", Tasks: data.discovery})
|
||||
if err != nil {
|
||||
return dataset{}, errors.New("encode SaaS test task discovery")
|
||||
}
|
||||
if err := contract.ValidateCurrent("task-discovery", page); err != nil {
|
||||
return dataset{}, errors.New("SaaS test task discovery violates the current contract")
|
||||
}
|
||||
return data, nil
|
||||
}
|
||||
|
||||
func (d dataset) handler(secret string) http.Handler {
|
||||
mux := http.NewServeMux()
|
||||
const prefix = "/internal/v1/dispatcher/"
|
||||
write := func(w http.ResponseWriter, body []byte) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, _ = w.Write(body)
|
||||
}
|
||||
mux.HandleFunc("GET "+prefix+"sip", func(w http.ResponseWriter, _ *http.Request) { write(w, d.sip) })
|
||||
mux.HandleFunc("GET "+prefix+"ai-providers", func(w http.ResponseWriter, _ *http.Request) { write(w, d.providers) })
|
||||
mux.HandleFunc("GET "+prefix+"task/{task_id}", func(w http.ResponseWriter, r *http.Request) {
|
||||
body, exists := d.tasks[r.PathValue("task_id")]
|
||||
if !exists {
|
||||
http.Error(w, "unknown task", http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
write(w, body)
|
||||
})
|
||||
mux.HandleFunc("GET "+prefix+"tenant/{tenant_id}/quota", func(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := strconv.ParseInt(r.PathValue("tenant_id"), 10, 64)
|
||||
if err != nil || id != d.tenantID {
|
||||
http.Error(w, "unknown tenant", http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
write(w, d.quota)
|
||||
})
|
||||
mux.HandleFunc("GET "+prefix+"tasks", func(w http.ResponseWriter, r *http.Request) {
|
||||
params := r.URL.Query()["after"]
|
||||
if len(params) > 1 || (len(params) == 1 && params[0] != "" && params[0] != "mock-complete") {
|
||||
http.Error(w, "unknown task discovery cursor", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
tasks := d.discovery
|
||||
if len(params) == 1 && params[0] == "mock-complete" {
|
||||
tasks = []configread.DiscoveredTask{}
|
||||
}
|
||||
body, err := json.Marshal(configread.TaskPage{DispatcherID: d.dispatcherID, Cursor: "mock-complete", Tasks: tasks})
|
||||
if err != nil {
|
||||
http.Error(w, "encode task discovery", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
write(w, body)
|
||||
})
|
||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if !strings.HasPrefix(r.URL.Path, prefix) {
|
||||
http.NotFound(w, r)
|
||||
return
|
||||
}
|
||||
if subtle.ConstantTimeCompare([]byte(r.Header.Get("X-DISPATCHER-id")), []byte(d.dispatcherID)) != 1 ||
|
||||
subtle.ConstantTimeCompare([]byte(r.Header.Get("X-DISPATCHER-SECRET-KEY")), []byte(secret)) != 1 {
|
||||
http.Error(w, "unknown dispatcher", http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
mux.ServeHTTP(w, r)
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user