Verify current HTTP reads and fail-closed outages

This commit is contained in:
2026-09-30 09:05:40 +08:00
parent 0e0537ebb5
commit 2b95966ba6
5 changed files with 148 additions and 17 deletions
@@ -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):
}
}
@@ -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)
+1 -1
View File
@@ -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
@@ -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 {
+53
View File
@@ -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) {