diff --git a/AGENTS.md b/AGENTS.md index 90f4764..dec9e75 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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 的现行边界 diff --git a/deploys/test/saas-mock/README.md b/deploys/test/saas-mock/README.md new file mode 100644 index 0000000..6134ada --- /dev/null +++ b/deploys/test/saas-mock/README.md @@ -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` 必须相同,任务须属于同一租户且文件名为 `.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 --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 或生产签收。 diff --git a/deploys/test/saas-mock/main.go b/deploys/test/saas-mock/main.go new file mode 100644 index 0000000..6a50617 --- /dev/null +++ b/deploys/test/saas-mock/main.go @@ -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)) +} diff --git a/deploys/test/saas-mock/main_test.go b/deploys/test/saas-mock/main_test.go new file mode 100644 index 0000000..13cfdbd --- /dev/null +++ b/deploys/test/saas-mock/main_test.go @@ -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") + } +} diff --git a/deploys/test/saas-mock/mq.go b/deploys/test/saas-mock/mq.go new file mode 100644 index 0000000..58ecc20 --- /dev/null +++ b/deploys/test/saas-mock/mq.go @@ -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 +} diff --git a/deploys/test/saas-mock/mq_test.go b/deploys/test/saas-mock/mq_test.go new file mode 100644 index 0000000..12bd0d9 --- /dev/null +++ b/deploys/test/saas-mock/mq_test.go @@ -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) + } +} diff --git a/deploys/test/saas-mock/server.go b/deploys/test/saas-mock/server.go new file mode 100644 index 0000000..027bc4b --- /dev/null +++ b/deploys/test/saas-mock/server.go @@ -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) + }) +}