From 4f319a58fa34332a9a7c0c8e5738606df78b4434 Mon Sep 17 00:00:00 2001 From: Rogee Date: Sun, 4 Oct 2026 15:57:41 +0800 Subject: [PATCH] feat(test): publish one authorized SaaS command over confirmed MQ --- AGENTS.md | 2 +- deploys/test/saas-mock/README.md | 5 +- deploys/test/saas-mock/main.go | 22 ++++- deploys/test/saas-mock/mq.go | 9 +- deploys/test/saas-mock/publish.go | 106 +++++++++++++++++++++++ deploys/test/saas-mock/publish_test.go | 112 +++++++++++++++++++++++++ scripts/check-current-mq-mock.sh | 6 +- 7 files changed, 256 insertions(+), 6 deletions(-) create mode 100644 deploys/test/saas-mock/publish.go create mode 100644 deploys/test/saas-mock/publish_test.go diff --git a/AGENTS.md b/AGENTS.md index ed32e6d..7119a27 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -81,7 +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、单租户**。根命令只接受显式 `agent`/`dispatcher`;`mixed` 和裸 `real` 仍拒绝;隔离 `mock`、只读 `sip-only` 和需完整非生产证据/正式获批指令的 `nonprod-real` 为彼此独立的入口。新增非生产真实入口尚未在测试机完成全链路核验或拨号;不得以编译/本机 Mock 通过宣称可用。另有严格隔离的 `--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。测试用正式入口必须同时具备只读配置、专用 MQ、真实 Agent/Asterisk/AI/录音/OSS、逐通确认和拨号前活跃抓包证据,缺一项不得试拨。此模拟不构成真实 SaaS 签收。 +- 使用者批准独立的测试专用 `deploys/test/saas-mock/` 仅模拟缺失的 SaaS:从 0600 的静态快照提供五类正式 HTTP 配置,在专用 RabbitMQ vhost 中由模拟 SaaS 预建现行拓扑;不在 Dispatcher 内注入快照;默认不发布外呼消息,只有受限脚本已为同一 `event_id`/线路/原始号码启动抓包后,才可显式单次 MQ 投递;不替代 Agent/Asterisk/AI/OSS。测试用正式入口必须同时具备只读配置、专用 MQ、真实 Agent/Asterisk/AI/录音/OSS、逐通确认和拨号前活跃抓包证据,缺一项不得试拨。此模拟不构成真实 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 index 6134ada..234e4c9 100644 --- a/deploys/test/saas-mock/README.md +++ b/deploys/test/saas-mock/README.md @@ -1,6 +1,6 @@ # SaaS 侧测试模拟服务(非生产) -仅替代当前缺失的 SaaS,不改 Dispatcher 的正式 HTTP 配置接口或 MQ 归属。**当前只提供静态配置和预建队列;不会投递 `call.execute`、拨号或产生业务结果。** 虚构的合同示例不能充当真实 AI、线路或任务授权。 +仅替代当前缺失的 SaaS,不改 Dispatcher 的正式 HTTP 配置接口或 MQ 归属。默认只提供静态配置和预建队列;另有**显式单次** `call.execute` 投递命令,绝不自动拨号或产生业务结果。 虚构的合同示例不能充当真实 AI、线路或任务授权。 ## 数据 @@ -12,7 +12,8 @@ - 准备测试 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 run ./deploys/test/saas-mock --data <私有目录> --dispatcher-id --publish-event-id <唯一事件号> --publish-task-id <单线路任务号> --publish-callee <原始白名单号码>`。必须在 `nonprod-call-evidence.sh --call-id <同一事件号> --trunk <任务唯一线路> --target <同一原始号码> -- <单次投递命令>` 启用并确认 SIP/RTP 抓包、PJSIP logger 和主机门禁之后运行;不得预投、批量投递、自动重试或换线。任务快照须只允许一条真实线路,投递仅含任务号和原始号码。RabbitMQ 必须用上述专用环境变量,不能借用默认或共享 vhost。发布确认只代表 MQ 接收,不代表 SaaS 已收到最终结果。 测试:`go test ./deploys/test/saas-mock` 验证正式配置客户端;设置指向**单独隔离 vhost** 的 `SAAS_MOCK_TEST_BROKER_URL` 后,`TestSaaSMockProvisionsDispatcherTopology` 还将实际预建 MQ 并用 Dispatcher 被动读回。缺省测试不会连接共享 RabbitMQ。 -当前正式呼叫仍只允许 Mock,测试机 ARI/HTTP、真实 AI/TTS、录音、OSS 与拨号前抓包尚未完成验证。六组外呼均未执行。此服务通过测试也不代表真实 SaaS 或生产签收。 +正式入口已经具备独立 `nonprod-real` 代码路径,但测试机的正式链路、真实 AI/TTS、录音、OSS 与拨号前抓包仍未完成整体验证。六组外呼均未执行。此服务通过测试也不代表真实 SaaS 或生产签收。 diff --git a/deploys/test/saas-mock/main.go b/deploys/test/saas-mock/main.go index 6a50617..604c87a 100644 --- a/deploys/test/saas-mock/main.go +++ b/deploys/test/saas-mock/main.go @@ -1,8 +1,9 @@ // saas-mock serves approved, static SaaS test snapshots through the formal -// Dispatcher configuration API. It never creates calls or publishes MQ jobs. +// Dispatcher configuration API. One-shot publishing is an explicit separate mode. package main import ( + "context" "errors" "flag" "log" @@ -17,9 +18,28 @@ func main() { 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") + publishID := flag.String("publish-event-id", "", "one-shot MQ event ID matching capture --call-id") + publishTask := flag.String("publish-task-id", "", "one approved single-trunk task") + publishCallee := flag.String("publish-callee", "", "one original allowlisted callee") flag.Parse() secret := os.Getenv("SAAS_MOCK_DISPATCHER_SECRET") brokerURL := os.Getenv("SAAS_MOCK_RABBITMQ_URL") + if *publishID != "" || *publishTask != "" || *publishCallee != "" { + if *dataDir == "" || *dispatcherID == "" || brokerURL == "" || *publishID == "" || *publishTask == "" || *publishCallee == "" || flag.NArg() != 0 { + log.Fatal("one-shot publish requires private data, dispatcher ID, MQ environment and explicit event/task/callee") + } + data, err := loadDataset(*dataDir, *dispatcherID) + if err != nil { + log.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := publishExecute(ctx, brokerURL, data, *publishID, *publishTask, *publishCallee); err != nil { + log.Fatal(err) + } + log.Printf("one-shot command broker-confirmed: event_id=%q task_id=%q (not SaaS result receipt)", *publishID, *publishTask) + return + } 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")) } diff --git a/deploys/test/saas-mock/mq.go b/deploys/test/saas-mock/mq.go index 58ecc20..2050521 100644 --- a/deploys/test/saas-mock/mq.go +++ b/deploys/test/saas-mock/mq.go @@ -12,12 +12,19 @@ import ( // 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 { +func requireDedicatedVhost(brokerURL string) 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") } + return nil +} + +func provisionMQ(brokerURL string, data dataset) error { + if err := requireDedicatedVhost(brokerURL); err != nil { + return err + } conn, err := amqp.Dial(brokerURL) if err != nil { return fmt.Errorf("connect SaaS test MQ: %T", err) diff --git a/deploys/test/saas-mock/publish.go b/deploys/test/saas-mock/publish.go new file mode 100644 index 0000000..2671523 --- /dev/null +++ b/deploys/test/saas-mock/publish.go @@ -0,0 +1,106 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "regexp" + "time" + + "git.ipao.vip/rogee/go-sip/internal/configread" + "git.ipao.vip/rogee/go-sip/internal/contract" + "git.ipao.vip/rogee/go-sip/internal/mq" + amqp "github.com/rabbitmq/amqp091-go" +) + +var mockCallID = regexp.MustCompile(`^[A-Za-z0-9._-]{1,100}$`) + +// buildExecute deliberately carries only the two approved call inputs. Trunk, +// caller, AI and duration remain immutable properties of the SaaS task read. +func buildExecute(data dataset, eventID, taskID, callee string, now time.Time) (string, []byte, error) { + shanghai, err := time.LoadLocation("Asia/Shanghai") + if err != nil { + return "", nil, err + } + hour := now.In(shanghai).Hour() + if !mockCallID.MatchString(eventID) || !mockCallID.MatchString(taskID) || + (callee != "15003164745" && callee != "15830461047") || hour < 9 || hour >= 20 { + return "", nil, errors.New("one-shot command identity, whitelist or real call window rejected") + } + body, ok := data.tasks[taskID] + if !ok { + return "", nil, errors.New("one-shot command task is absent from the approved SaaS dataset") + } + var task configread.Task + if err := json.Unmarshal(body, &task); err != nil || task.DispatcherID != data.dispatcherID || task.TenantID != data.tenantID || task.TaskID != taskID || task.Status != "running" || len(task.AllowedTrunkIDs) != 1 { + return "", nil, errors.New("one-shot command requires a running task pinned to exactly one approved trunk") + } + command, err := json.Marshal(struct { + EventID string `json:"event_id"` + EventType string `json:"event_type"` + DispatcherID string `json:"dispatcher_id"` + TenantID int64 `json:"tenant_id"` + IssuedAt string `json:"issued_at"` + Payload struct { + TaskID string `json:"task_id"` + Callee string `json:"callee"` + } `json:"payload"` + }{EventID: eventID, EventType: "call.execute", DispatcherID: data.dispatcherID, TenantID: data.tenantID, IssuedAt: now.UTC().Format(time.RFC3339Nano), Payload: struct { + TaskID string `json:"task_id"` + Callee string `json:"callee"` + }{TaskID: taskID, Callee: callee}}) + if err != nil || contract.ValidateCurrent("mq", command) != nil { + return "", nil, errors.New("one-shot command violates current MQ contract") + } + return "d." + data.dispatcherID + ".task." + taskID + ".in", command, nil +} + +// publishExecute publishes exactly once with mandatory routing and publisher +// confirmation. An ambiguous response is NOT retried or called SaaS receipt. +func publishExecute(ctx context.Context, brokerURL string, data dataset, eventID, taskID, callee string) error { + return publishExecuteAt(ctx, brokerURL, data, eventID, taskID, callee, time.Now()) +} + +func publishExecuteAt(ctx context.Context, brokerURL string, data dataset, eventID, taskID, callee string, now time.Time) error { + if err := requireDedicatedVhost(brokerURL); err != nil { + return err + } + routingKey, body, err := buildExecute(data, eventID, taskID, callee, now) + if err != nil { + return err + } + conn, err := amqp.Dial(brokerURL) + if err != nil { + return fmt.Errorf("connect one-shot test broker: %T", err) + } + defer conn.Close() + ch, err := conn.Channel() + if err != nil { + return fmt.Errorf("open one-shot publish channel: %T", err) + } + defer ch.Close() + if err := ch.Confirm(false); err != nil { + return fmt.Errorf("enable one-shot publisher confirmation: %T", err) + } + returned := ch.NotifyReturn(make(chan amqp.Return, 1)) + confirmed := ch.NotifyPublish(make(chan amqp.Confirmation, 1)) + if err := ch.PublishWithContext(ctx, mq.CommandsExchange, routingKey, true, false, + amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil { + return fmt.Errorf("one-shot publish outcome unknown: %T", err) + } + select { + case <-ctx.Done(): + return errors.New("one-shot publish confirmation unknown; do not republish") + case confirmation, ok := <-confirmed: + if !ok || !confirmation.Ack { + return errors.New("one-shot command was not publisher-confirmed; do not republish") + } + select { + case <-returned: + return errors.New("one-shot command was unroutable") + default: + return nil // broker accepted the command, NOT SaaS application receipt + } + } +} diff --git a/deploys/test/saas-mock/publish_test.go b/deploys/test/saas-mock/publish_test.go new file mode 100644 index 0000000..b082c22 --- /dev/null +++ b/deploys/test/saas-mock/publish_test.go @@ -0,0 +1,112 @@ +package main + +import ( + "context" + "encoding/json" + "os" + "strings" + "testing" + "time" + + "git.ipao.vip/rogee/go-sip/internal/contract" + "github.com/google/uuid" + amqp "github.com/rabbitmq/amqp091-go" +) + +func TestSaaSMockPublishesOneCommandOnProvisionedRabbitMQ(t *testing.T) { + url := os.Getenv("SAAS_MOCK_TEST_BROKER_URL") + if url == "" { + t.Skip("requires an isolated ephemeral dedicated-vhost RabbitMQ") + } + data, err := loadDataset(testDataDir(t), testDispatcher) + if err != nil { + t.Fatal(err) + } + if err := provisionMQ(url, data); err != nil { + t.Fatal(err) + } + found := false + for _, discovered := range data.discovery { + found = found || discovered.TaskID == "task-full" + } + if !found { + t.Fatal("fixture task-full was not preprovisioned") + } + conn, err := amqp.Dial(url) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + ch, err := conn.Channel() + if err != nil { + t.Fatal(err) + } + defer ch.Close() + eventID := "test-" + uuid.NewString() + at := time.Date(2026, 10, 4, 2, 0, 0, 0, time.UTC) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := publishExecuteAt(ctx, url, data, eventID, "task-full", "15003164745", at); err != nil { + t.Fatal(err) + } + queue := "agent-call.d." + testDispatcher + ".task.task-full.v1" + message, ok, err := ch.Get(queue, false) + if err != nil || !ok { + t.Fatalf("one command must reach exactly its preprovisioned task queue: present=%v err=%v", ok, err) + } + if err := contract.ValidateCurrent("mq", message.Body); err != nil || !strings.Contains(string(message.Body), eventID) || message.DeliveryMode != amqp.Persistent { + t.Fatalf("broker delivered wrong command identity or persistence: err=%v", err) + } + if err := message.Ack(false); err != nil { + t.Fatal(err) + } + if extra, ok, err := ch.Get(queue, false); err != nil || ok { + if ok { + _ = extra.Nack(false, true) + } + t.Fatalf("one publish created a second command: present=%v err=%v", ok, err) + } +} + +func TestSaaSMockBuildsOneApprovedCommandForTheBoundTask(t *testing.T) { + data, err := loadDataset(testDataDir(t), testDispatcher) + if err != nil { + t.Fatal(err) + } + inside := time.Date(2026, 10, 4, 2, 0, 0, 0, time.UTC) // Shanghai 10:00 + key, body, err := buildExecute(data, "event-once-1", "task-full", "15003164745", inside) + if err != nil || key != "d."+testDispatcher+".task.task-full.in" { + t.Fatalf("SaaS command must target only the assigned Dispatcher/task: key=%q err=%v", key, err) + } + if err := contract.ValidateCurrent("mq", body); err != nil { + t.Fatalf("call command violates the formal contract: %v", err) + } + var event struct { + EventID string `json:"event_id"` + Type string `json:"event_type"` + DispatcherID string `json:"dispatcher_id"` + TenantID int64 `json:"tenant_id"` + Payload struct { + TaskID string `json:"task_id"` + Callee string `json:"callee"` + } `json:"payload"` + } + if err := json.Unmarshal(body, &event); err != nil || event.EventID != "event-once-1" || event.Type != "call.execute" || event.DispatcherID != testDispatcher || event.TenantID != 1001 || event.Payload.TaskID != "task-full" || event.Payload.Callee != "15003164745" || strings.Contains(string(body), "trunk-mock") { + t.Fatalf("SaaS must not leak a trunk/caller/AI override into the execute command: %+v err=%v", event, err) + } + for _, test := range []struct { + id, task, callee string + at time.Time + }{ + {"../event", "task-full", "15003164745", inside}, + {"event-2", "missing", "15003164745", inside}, + {"event-3", "task-full", "708915003164745", inside}, + {"event-4", "task-full", "13900000000", inside}, + {"event-5", "task-full", "15003164745", inside.Add(-2 * time.Hour)}, + {"event-6", "task-full", "15003164745", inside.Add(10 * time.Hour)}, + } { + if _, _, err := buildExecute(data, test.id, test.task, test.callee, test.at); err == nil { + t.Fatalf("invalid or out-of-window command was allowed: event=%q task=%q", test.id, test.task) + } + } +} diff --git a/scripts/check-current-mq-mock.sh b/scripts/check-current-mq-mock.sh index ed6fe9b..9f56401 100644 --- a/scripts/check-current-mq-mock.sh +++ b/scripts/check-current-mq-mock.sh @@ -45,12 +45,15 @@ d_pw=$(openssl rand -hex 12) if ! docker exec -u rabbitmq "$name" rabbitmqctl add_user saas_mock "$saas_pw" >/dev/null 2>&1 || ! docker exec -u rabbitmq "$name" rabbitmqctl add_user dispatcher_mock "$d_pw" >/dev/null 2>&1 || ! docker exec -u rabbitmq "$name" rabbitmqctl set_permissions -p / saas_mock '.*' '.*' '.*' >/dev/null 2>&1 || - ! docker exec -u rabbitmq "$name" rabbitmqctl set_permissions -p / dispatcher_mock '^$' '^agent-call\.saas\.v1$' '^agent-call\.d\..*\.(control|task\..*)\.v1$' >/dev/null 2>&1; then + ! docker exec -u rabbitmq "$name" rabbitmqctl set_permissions -p / dispatcher_mock '^$' '^agent-call\.saas\.v1$' '^agent-call\.d\..*\.(control|task\..*)\.v1$' >/dev/null 2>&1 || + ! docker exec -u rabbitmq "$name" rabbitmqctl add_vhost saas-mock-go-sip-test >/dev/null 2>&1 || + ! docker exec -u rabbitmq "$name" rabbitmqctl set_permissions -p saas-mock-go-sip-test saas_mock '.*' '.*' '.*' >/dev/null 2>&1; then echo 'isolated RabbitMQ Mock user/permission setup failed' >&2 exit 1 fi export RABBITMQ_URL="amqp://dispatcher_mock:${d_pw}@127.0.0.1:${port}/" export RABBITMQ_PROVISIONER_URL="amqp://saas_mock:${saas_pw}@127.0.0.1:${port}/" +export SAAS_MOCK_TEST_BROKER_URL="amqp://saas_mock:${saas_pw}@127.0.0.1:${port}/saas-mock-go-sip-test" run_required_test() { local package=$1 test_name=$2 output if ! output=$(go test -tags=integration "$package" -run "^${test_name}$" -count=1 -v 2>&1); then @@ -66,3 +69,4 @@ run_required_test() { run_required_test ./internal/mq TestBrokerSharedResultQueueAndNoConfigure run_required_test ./internal/dispatcher TestRuntimeIsolatedControlBacklogExecuteAndSharedResult run_required_test ./cmd/sip-go-agent TestDispatcherCommandStartsWithIsolatedMQHTTPAndAgent +run_required_test ./deploys/test/saas-mock TestSaaSMockPublishesOneCommandOnProvisionedRabbitMQ