diff --git a/cmd/sip-go-agent/current_dispatcher_integration_test.go b/cmd/sip-go-agent/current_dispatcher_integration_test.go index 76a2412..d8f8e19 100644 --- a/cmd/sip-go-agent/current_dispatcher_integration_test.go +++ b/cmd/sip-go-agent/current_dispatcher_integration_test.go @@ -71,7 +71,18 @@ func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) t.Fatal(err) } defer func() { _ = admin.QueueUnbind(result.Queue, result.BindingKey, result.Exchange, nil) }() - for _, queue := range []string{control.Queue, result.Queue} { + taskRoute, err := tenant.CurrentTaskRoute(dispatcherID, "task-asr") + if err != nil { + t.Fatal(err) + } + if _, err := admin.QueueDeclare(taskRoute.Queue, true, false, false, false, nil); err != nil { + t.Fatal(err) + } + defer func() { _, _ = admin.QueueDelete(taskRoute.Queue, false, false, false) }() + if err := admin.QueueBind(taskRoute.Queue, taskRoute.BindingKey, taskRoute.Exchange, false, nil); err != nil { + t.Fatal(err) + } + for _, queue := range []string{control.Queue, result.Queue, taskRoute.Queue} { if _, err := admin.QueuePurge(queue, false); err != nil { t.Fatal(err) } @@ -106,15 +117,15 @@ func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) go func() { _ = agentGRPC.Serve(agentListener) }() defer agentGRPC.Stop() - sipBody, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", "config-read-sip.json")) - if err != nil { - t.Fatal(err) + examples := map[string][]byte{} + for _, name := range []string{"config-read-sip.json", "task-discovery-page.json", "task-discovery-end.json", "config-read-task-asr.json", "config-read-providers.json", "config-read-quota.json"} { + body, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", name)) + if err != nil { + t.Fatal(err) + } + examples[name] = body } - taskBody, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", "task-discovery-end.json")) - if err != nil { - t.Fatal(err) - } - var sipReads, taskReads atomic.Int64 + var sipReads, discoveryReads, taskReads, providerReads, quotaReads atomic.Int64 saas := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if r.Header.Get("X-DISPATCHER-id") != dispatcherID || r.Header.Get("X-DISPATCHER-SECRET-KEY") != "isolated-secret" { http.Error(w, "unapproved Dispatcher", http.StatusForbidden) @@ -124,10 +135,26 @@ func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) switch r.URL.Path { case "/internal/v1/dispatcher/sip": sipReads.Add(1) - _, _ = w.Write(sipBody) + _, _ = w.Write(examples["config-read-sip.json"]) case "/internal/v1/dispatcher/tasks": + discoveryReads.Add(1) + switch r.URL.Query().Get("after") { + case "": + _, _ = w.Write(examples["task-discovery-page.json"]) + case "opaque-page-token-1", "opaque-end-token": + _, _ = w.Write(examples["task-discovery-end.json"]) + default: + http.Error(w, "unexpected Mock cursor", http.StatusBadRequest) + } + case "/internal/v1/dispatcher/task/task-asr": taskReads.Add(1) - _, _ = w.Write(taskBody) + _, _ = w.Write(examples["config-read-task-asr.json"]) + case "/internal/v1/dispatcher/ai-providers": + providerReads.Add(1) + _, _ = w.Write(examples["config-read-providers.json"]) + case "/internal/v1/dispatcher/tenant/1001/quota": + quotaReads.Add(1) + _, _ = w.Write(examples["config-read-quota.json"]) default: http.Error(w, "unexpected Mock read", http.StatusNotFound) } @@ -176,14 +203,33 @@ func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) finished := make(chan error, 1) go func() { finished <- command.Execute() }() deadline := time.After(8 * time.Second) - for sipReads.Load() == 0 || taskReads.Load() == 0 { + for sipReads.Load() < 2 || discoveryReads.Load() < 2 || taskReads.Load() == 0 || providerReads.Load() == 0 || quotaReads.Load() == 0 { select { case err := <-finished: cancel() - t.Fatalf("isolated Dispatcher exited before HTTP bootstrap: %v", err) + t.Fatalf("isolated Dispatcher exited before the five HTTP resources were read: %v", err) case <-deadline: cancel() - t.Fatal("isolated Dispatcher never read SIP and the empty task snapshot") + t.Fatal("isolated Dispatcher did not read SIP, discovery, task, providers, and tenant quota") + case <-time.After(20 * time.Millisecond): + } + } + for { + queue, err := admin.QueueInspect(taskRoute.Queue) + if err != nil { + cancel() + t.Fatalf("assigned task queue cannot be inspected: %v", err) + } + if queue.Consumers > 0 { + break + } + select { + case err := <-finished: + cancel() + t.Fatalf("isolated Dispatcher exited before consuming its task queue: %v", err) + case <-deadline: + cancel() + t.Fatal("isolated Dispatcher did not consume its SaaS-provisioned task queue") case <-time.After(20 * time.Millisecond): } } diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index ae5da9f..01e37b0 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -45,10 +45,11 @@ ## P03:HTTP 读取分批改造(未整体签收) -- `contract.ValidateCurrent` 与 `configread` 按当前 Schema 读取 SIP、provider、task、quota 和 cursor 任务发现;严格检查数字 tenant_id、本 D 归属及不可变配置。provider 凭据原值只保留在内存快照,不写日志;Agent 参数中的显式 0/false 保真;无旧 Schema/旧配置回退。 +- `contract.ValidateCurrent` 与 `configread` 按当前 Schema 读取 SIP、provider、task、quota 和 cursor 任务发现;严格检查数字 tenant_id、本 D 归属及不可变配置。provider 凭据原值保留在持久快照并经受信 Agent RPC 交付,不写入日志或证据;Agent 参数中的显式 0/false 保真;无旧 Schema/旧配置回退。 - CLI 组装辅助 `dispatcherConfigurationClient` 将严格 Mock 环境预检与当前 HTTP Client 绑定;本机 HTTP 隔离测试实际读取 SIP,核对固定 `/internal/v1/dispatcher/sip`、D 身份/密钥 Header、revision 与不打开 SQLite;旧 `DISPATCHER_CONFIG_READ_BASE_URL` 不可充当缺失的当前地址。`go test ./cmd/sip-go-agent -run '^TestDispatcherConfigurationClient' -count=1` 通过;主 `dispatcher` 命令已通过该辅助读取本机 Mock;真实 SaaS 的授权和连通尚未验证。 - `store.OpenCurrent` 新建数字租户 SQLite 状态;旧表、旧版当前布局、残缺布局均在写入前拒绝并保留原记录;不执行旧数据迁移或自动清理。启动时完整发现同一快照一次提交,分页增量逐页持久提交后才推进**内存** cursor;失败关闭准入,重启重新取完整快照。HTTP 的旧 running 不能解除 MQ 暂停/终止,同 revision 异内容及跨任务 SIP/租户额度冲突拒绝。 -- `CurrentBootstrap` 先关闭准入,核验 SIP 全量与 Agent/Asterisk 已加载 revision、读取任务和 provider/额度,再排空 MQ 控制积压,最后依据已验证 SIP revision 开准入;有更新的持久 SIP 通知时保持关闭但控制与结果处理仍可继续。`CurrentDiscoveryFollower` 逐页绑定任务快照;HTTP 错误、失效或授权不一致只失败,不回退旧读取。**根 `dispatcher` 命令已接入该启动链路,空任务隔离快照通过;含任务的主进程 provider 交付与真实 Agent/Asterisk 加载仍未验证。** +- `CurrentBootstrap` 先关闭准入,核验 SIP 全量与 Agent/Asterisk 已加载 revision、读取任务和 provider/额度,再排空 MQ 控制积压,最后依据已验证 SIP revision 开准入;有更新的持久 SIP 通知时保持关闭但控制与结果处理仍可继续。`CurrentDiscoveryFollower` 逐页绑定任务快照;HTTP 错误、失效或授权不一致只失败,不回退旧读取。**根 `dispatcher` 命令已接入该启动链路;本机隔离测试实际读取 SIP、分页发现、任务、provider 与租户额度,核验模拟 Agent 已加载 revision 后订阅 SaaS 预建任务队列。主进程实际呼叫、录音与结果交付,以及真实 Agent/Asterisk 加载仍未验证。** +- HTTP 失败回归覆盖 SIP、provider、task、tenant quota 四个读取端点:已有成功响应后收到 503 必须失败且保留 HTTP 状态与资源归属,不复用旧响应或泄漏受控 Header。任务发现同样在后续 503 时拒绝旧页;`TestReadCurrentTasksFailsOnUnavailableRefreshWithoutStaleFallback` 先暴露缺少错误归属,修复后保留 `HTTPError` 及 `task discovery` 上下文。根命令的预建队列、完整 HTTP 读取与模拟 Agent 激活由 `bash scripts/check-current-mq-mock.sh` 隔离验证。 - TDD 与回归:`go test ./internal/configread ./internal/tenant ./internal/store ./internal/dispatcher -count=1`、`bash scripts/check-current-contracts.sh`、已提交 `a0118e3` 的干净归档测试通过;旧布局行/表原样保留由 `TestCurrentStoreRefusesPreviousCurrentLayoutBeforeModifyingDatabase` 覆盖。 ## P04:隔离 MQ、控制与外呼接纳(仅项目内 Mock) diff --git a/internal/configread/current_discovery.go b/internal/configread/current_discovery.go index 5c54c1e..94a3232 100644 --- a/internal/configread/current_discovery.go +++ b/internal/configread/current_discovery.go @@ -34,7 +34,7 @@ func (c *Client) ReadCurrentTasks(ctx context.Context, after string) (CurrentTas } body, err := c.get(ctx, path) if err != nil { - return CurrentTaskPage{}, err + return CurrentTaskPage{}, fmt.Errorf("read task discovery: %w", err) } if err := contract.ValidateCurrent("task-discovery", body); err != nil { return CurrentTaskPage{}, err diff --git a/internal/configread/current_discovery_test.go b/internal/configread/current_discovery_test.go index ada52bb..4512b0b 100644 --- a/internal/configread/current_discovery_test.go +++ b/internal/configread/current_discovery_test.go @@ -2,9 +2,11 @@ package configread import ( "context" + "errors" "net/http" "net/http/httptest" "strings" + "sync/atomic" "testing" ) @@ -42,6 +44,35 @@ func TestReadAllCurrentTasksContinuesAfterShortPageUntilEmptyPage(t *testing.T) } } +func TestReadCurrentTasksFailsOnUnavailableRefreshWithoutStaleFallback(t *testing.T) { + var unavailable atomic.Bool + body := currentExample(t, "task-discovery-page") + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + if unavailable.Load() { + w.WriteHeader(http.StatusServiceUnavailable) + _, _ = w.Write([]byte(`{"resource":"error","error":{"code":"unavailable","message":"try later"}}`)) + return + } + _, _ = w.Write(body) + })) + defer server.Close() + client, err := NewClient(server.URL, "c046b893-8628-4589-ae50-619d049248a6", "test-secret", server.Client()) + if err != nil { + t.Fatal(err) + } + page, err := client.ReadCurrentTasks(context.Background(), "") + if err != nil || len(page.Tasks) != 1 { + t.Fatalf("expected a fresh assigned task before outage: page=%+v err=%v", page, err) + } + unavailable.Store(true) + _, err = client.ReadCurrentTasks(context.Background(), page.Cursor) + var httpErr *HTTPError + if !errors.As(err, &httpErr) || httpErr.StatusCode != http.StatusServiceUnavailable || !strings.Contains(err.Error(), "task discovery") { + t.Fatalf("refresh failure lost its owner/status or reused stale data: %v", err) + } +} + func TestReadCurrentTasksFailClosed(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" for _, tc := range []struct { diff --git a/internal/configread/current_test.go b/internal/configread/current_test.go index bfab19a..644db42 100644 --- a/internal/configread/current_test.go +++ b/internal/configread/current_test.go @@ -3,11 +3,13 @@ package configread import ( "context" "encoding/json" + "errors" "net/http" "net/http/httptest" "os" "path/filepath" "strings" + "sync/atomic" "testing" ) @@ -103,6 +105,57 @@ func TestReadCurrentSIPWithoutTasks(t *testing.T) { } } +func TestReadCurrentTaskDoesNotReuseConfigAfterHTTPFailure(t *testing.T) { + const dispatcherID = "c046b893-8628-4589-ae50-619d049248a6" + responses := map[string][]byte{ + "/internal/v1/dispatcher/sip": currentExample(t, "config-read-sip"), + "/internal/v1/dispatcher/ai-providers": currentExample(t, "config-read-providers"), + "/internal/v1/dispatcher/task/task-asr": currentExample(t, "config-read-task-asr"), + "/internal/v1/dispatcher/tenant/1001/quota": currentExample(t, "config-read-quota"), + } + for _, tc := range []struct{ path, owner string }{ + {"/internal/v1/dispatcher/sip", "SIP"}, + {"/internal/v1/dispatcher/ai-providers", "AI providers"}, + {"/internal/v1/dispatcher/task/task-asr", "task configuration"}, + {"/internal/v1/dispatcher/tenant/1001/quota", "tenant quota"}, + } { + t.Run(tc.owner, func(t *testing.T) { + var unavailable atomic.Bool + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + if unavailable.Load() && r.URL.Path == tc.path { + w.WriteHeader(http.StatusServiceUnavailable) + _, _ = w.Write([]byte(`{"resource":"error","error":{"code":"unavailable","message":"try later"}}`)) + return + } + body, ok := responses[r.URL.Path] + if !ok { + http.NotFound(w, r) + return + } + _, _ = w.Write(body) + })) + defer server.Close() + client, err := NewClient(server.URL, dispatcherID, "isolated-secret", server.Client()) + if err != nil { + t.Fatal(err) + } + if _, err := client.ReadCurrentTask(context.Background(), "task-asr", 1001); err != nil { + t.Fatalf("approved resources failed before the outage: %v", err) + } + unavailable.Store(true) + _, err = client.ReadCurrentTask(context.Background(), "task-asr", 1001) + var httpErr *HTTPError + if !errors.As(err, &httpErr) || httpErr.StatusCode != http.StatusServiceUnavailable { + t.Fatalf("unavailable %s reused cached configuration or lost status: %v", tc.owner, err) + } + if !strings.Contains(strings.ToLower(err.Error()), strings.ToLower(tc.owner)) || strings.Contains(err.Error(), "isolated-secret") { + t.Fatalf("HTTP failure lost resource ownership or leaked credentials: %v", err) + } + }) + } +} + func TestReadCurrentTaskRejectsMismatchedOwnerAndMalformedSIP(t *testing.T) { for _, mode := range []string{"owner", "sip-revision", "provider-ref", "provider-disabled", "provider-wrong-role"} { t.Run(mode, func(t *testing.T) {