From 99a8c77132859993f5e9e067cce54e9fb620789f Mon Sep 17 00:00:00 2001 From: Rogee Date: Sun, 4 Oct 2026 17:11:57 +0800 Subject: [PATCH] feat(test): await and persist one confirmed SaaS call result --- deploys/test/saas-mock/README.md | 2 +- deploys/test/saas-mock/main.go | 18 ++- deploys/test/saas-mock/publish_test.go | 151 ++++++++++++++++++++ deploys/test/saas-mock/result_wait.go | 190 +++++++++++++++++++++++++ scripts/check-current-mq-mock.sh | 1 + 5 files changed, 359 insertions(+), 3 deletions(-) create mode 100644 deploys/test/saas-mock/result_wait.go diff --git a/deploys/test/saas-mock/README.md b/deploys/test/saas-mock/README.md index 3aed5d0..b0d0b66 100644 --- a/deploys/test/saas-mock/README.md +++ b/deploys/test/saas-mock/README.md @@ -13,7 +13,7 @@ - 准备测试 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 run ./deploys/test/saas-mock --data <私有目录> --dispatcher-id --publish-event-id <唯一事件号> --publish-task-id <单线路任务号> --publish-callee <原始白名单号码> --await-result-file <私有结果文件>`。此命令在专用结果队列中等待精确匹配的单通最终结果,先将原始结果写入 `0600` 私有文件并同步磁盘,才确认 MQ 消费;标准输出只显示结果摘要/hash,不输出转写、录音或签名 URL。必须在 `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。 diff --git a/deploys/test/saas-mock/main.go b/deploys/test/saas-mock/main.go index 0ab92f4..54b5a2b 100644 --- a/deploys/test/saas-mock/main.go +++ b/deploys/test/saas-mock/main.go @@ -22,11 +22,12 @@ func main() { 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") + awaitResultFile := flag.String("await-result-file", "", "absolute private path for durable final result; blocks until the exact call ends") flag.Parse() secret := os.Getenv("SAAS_MOCK_DISPATCHER_SECRET") brokerURL := os.Getenv("SAAS_MOCK_RABBITMQ_URL") if *validateOnly { - if *dataDir == "" || *dispatcherID == "" || *publishID != "" || *publishTask != "" || *publishCallee != "" || flag.NArg() != 0 { + if *dataDir == "" || *dispatcherID == "" || *publishID != "" || *publishTask != "" || *publishCallee != "" || *awaitResultFile != "" || flag.NArg() != 0 { log.Fatal("offline validation requires only a private data directory and bound Dispatcher ID") } data, err := loadDataset(*dataDir, *dispatcherID) @@ -36,7 +37,7 @@ func main() { log.Printf("private SaaS snapshots validated: task_count=%d", len(data.tasks)) return } - if *publishID != "" || *publishTask != "" || *publishCallee != "" { + if *publishID != "" || *publishTask != "" || *publishCallee != "" || *awaitResultFile != "" { 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") } @@ -44,6 +45,19 @@ func main() { if err != nil { log.Fatal(err) } + if *awaitResultFile != "" { + ctx, cancel := context.WithTimeout(context.Background(), 4*time.Minute) + defer cancel() + result, err := publishAndAwait(ctx, brokerURL, data, *publishID, *publishTask, *publishCallee, *awaitResultFile, time.Now()) + if err != nil { + log.Fatal(err) + } + log.Printf("one-shot SaaS result saved and acknowledged: event_id=%q outcome=%q recording_status=%q final_user_turns=%d result_sha256=%s", *publishID, result.Outcome, result.RecordingStatus, result.Transcripts, result.BodySHA256) + if result.Outcome != "answered" || result.RecordingStatus != "uploaded" { + log.Fatal("real call did not produce the required answered, uploaded-recording result; stop further attempts") + } + return + } ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() if err := publishExecute(ctx, brokerURL, data, *publishID, *publishTask, *publishCallee); err != nil { diff --git a/deploys/test/saas-mock/publish_test.go b/deploys/test/saas-mock/publish_test.go index b082c22..0f46eb5 100644 --- a/deploys/test/saas-mock/publish_test.go +++ b/deploys/test/saas-mock/publish_test.go @@ -1,9 +1,12 @@ package main import ( + "bytes" "context" "encoding/json" + "fmt" "os" + "path/filepath" "strings" "testing" "time" @@ -13,6 +16,154 @@ import ( amqp "github.com/rabbitmq/amqp091-go" ) +func TestOneShotResultMustMatchSourceIdentityAndPersistBeforeReceipt(t *testing.T) { + data, err := loadDataset(testDataDir(t), testDispatcher) + if err != nil { + t.Fatal(err) + } + sourceID := "event-once-1" + resultID := expectedResultEventID(data.dispatcherID, sourceID) + result := map[string]any{} + example, err := os.ReadFile("../../../contracts/local/examples/mq-result-uploaded.json") + if err != nil { + t.Fatal(err) + } + if err := json.Unmarshal(example, &result); err != nil { + t.Fatal(err) + } + result["event_id"] = resultID + result["dispatcher_id"] = data.dispatcherID + payload := result["payload"].(map[string]any) + payload["task_id"], payload["callee"], payload["trunk_id"] = "task-full", "15003164745", "trunk-mock" + body, err := json.Marshal(result) + if err != nil { + t.Fatal(err) + } + if err := contract.ValidateCurrent("mq", body); err != nil { + t.Fatal(err) + } + if _, err := summarizeOneShotResult(body, data.dispatcherID, sourceID, "task-full", "15003164745", "trunk-mock"); err != nil { + t.Fatal(err) + } + if _, err := summarizeOneShotResult(body, data.dispatcherID, "different-event", "task-full", "15003164745", "trunk-mock"); err == nil { + t.Fatal("unrelated final result must not be acknowledged") + } + if _, err := summarizeOneShotResult(body, data.dispatcherID, sourceID, "task-full", "15003164745", "different-trunk"); err == nil { + t.Fatal("result from the wrong trunk must not be acknowledged") + } + dir := t.TempDir() + if err := os.Chmod(dir, 0700); err != nil { + t.Fatal(err) + } + file := filepath.Join(dir, "one-result.json") + if err := persistOneShotResult(file, body); err != nil { + t.Fatal(err) + } + stored, err := os.ReadFile(file) + if err != nil || !bytes.Equal(stored, body) { + t.Fatal("final result not durably saved") + } + if info, err := os.Stat(file); err != nil || info.Mode().Perm() != 0600 { + t.Fatal("result must remain private") + } + if err := persistOneShotResult(file, body); err == nil { + t.Fatal("existing final result must not be overwritten") + } +} + +func TestSaaSMockAwaitsExactResultBeforeStoppingCapture(t *testing.T) { + url := os.Getenv("SAAS_MOCK_TEST_BROKER_URL") + if url == "" { + t.Skip("requires 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) + } + eventID := "test-" + uuid.NewString() + result, err := os.ReadFile("../../../contracts/local/examples/mq-result-uploaded.json") + if err != nil { + t.Fatal(err) + } + var message map[string]any + if err := json.Unmarshal(result, &message); err != nil { + t.Fatal(err) + } + message["event_id"] = expectedResultEventID(data.dispatcherID, eventID) + message["dispatcher_id"] = data.dispatcherID + payload := message["payload"].(map[string]any) + payload["task_id"], payload["callee"], payload["trunk_id"] = "task-full", "15003164745", "trunk-mock" + result, err = json.Marshal(message) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + producer := make(chan error, 1) + go func() { + conn, err := amqp.Dial(url) + if err != nil { + producer <- err + return + } + defer conn.Close() + ch, err := conn.Channel() + if err != nil { + producer <- err + return + } + defer ch.Close() + queue := "agent-call.d." + testDispatcher + ".task.task-full.v1" + for { + command, found, err := ch.Get(queue, false) + if err != nil { + producer <- err + return + } + if found { + if !bytes.Contains(command.Body, []byte(eventID)) { + producer <- fmt.Errorf("wrong command received") + return + } + if err := command.Ack(false); err != nil { + producer <- err + return + } + err = ch.PublishWithContext(ctx, "agent-call.saas.v1", "d."+testDispatcher+".out", true, false, amqp.Publishing{DeliveryMode: amqp.Persistent, ContentType: "application/json", Body: result}) + producer <- err + return + } + select { + case <-ctx.Done(): + producer <- ctx.Err() + return + case <-time.After(20 * time.Millisecond): + } + } + }() + dir := t.TempDir() + if err := os.Chmod(dir, 0700); err != nil { + t.Fatal(err) + } + file := filepath.Join(dir, "final.json") + observed, err := publishAndAwait(ctx, url, data, eventID, "task-full", "15003164745", file, time.Date(2026, 10, 4, 2, 0, 0, 0, time.UTC)) + if err != nil { + t.Fatal(err) + } + if err := <-producer; err != nil { + t.Fatal(err) + } + if observed.Outcome != "answered" || observed.RecordingStatus != "uploaded" { + t.Fatalf("result not acknowledged with real-looking recording: %+v", observed) + } + if saved, err := os.ReadFile(file); err != nil || !bytes.Equal(saved, result) { + t.Fatalf("result was not saved before broker receipt: err=%v", err) + } +} + func TestSaaSMockPublishesOneCommandOnProvisionedRabbitMQ(t *testing.T) { url := os.Getenv("SAAS_MOCK_TEST_BROKER_URL") if url == "" { diff --git a/deploys/test/saas-mock/result_wait.go b/deploys/test/saas-mock/result_wait.go new file mode 100644 index 0000000..0bb4ebe --- /dev/null +++ b/deploys/test/saas-mock/result_wait.go @@ -0,0 +1,190 @@ +package main + +import ( + "context" + "crypto/sha256" + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "time" + + "git.ipao.vip/rogee/go-sip/internal/configread" + "git.ipao.vip/rogee/go-sip/internal/contract" + amqp "github.com/rabbitmq/amqp091-go" +) + +const resultQueue = "agent-call.saas.events.v1" + +type oneShotResult struct { + Outcome string + RecordingStatus string + Transcripts int + BodySHA256 string +} + +func expectedResultEventID(dispatcherID, sourceID string) string { + identity := sha256.Sum256([]byte("call.execute.result\x00" + dispatcherID + "\x00" + sourceID)) + return fmt.Sprintf("result-%x", identity[:]) +} + +func summarizeOneShotResult(body []byte, dispatcherID, sourceID, taskID, callee, trunkID string) (oneShotResult, error) { + if err := contract.ValidateCurrent("mq", body); err != nil { + return oneShotResult{}, errors.New("final MQ result violates current contract") + } + var event struct { + EventID string `json:"event_id"` + EventType string `json:"event_type"` + DispatcherID string `json:"dispatcher_id"` + Payload struct { + TaskID string `json:"task_id"` + Callee string `json:"callee"` + TrunkID string `json:"trunk_id"` + Outcome string `json:"outcome"` + Transcript []json.RawMessage `json:"transcript"` + Recording struct { + Status string `json:"status"` + } `json:"recording"` + } `json:"payload"` + } + if err := json.Unmarshal(body, &event); err != nil || event.EventType != "call.execute.result" || event.EventID != expectedResultEventID(dispatcherID, sourceID) || event.DispatcherID != dispatcherID || event.Payload.TaskID != taskID || event.Payload.Callee != callee || event.Payload.TrunkID != trunkID { + return oneShotResult{}, errors.New("final MQ result does not belong to the selected one-shot call") + } + sha := sha256.Sum256(body) + return oneShotResult{Outcome: event.Payload.Outcome, RecordingStatus: event.Payload.Recording.Status, Transcripts: len(event.Payload.Transcript), BodySHA256: fmt.Sprintf("%x", sha[:])}, nil +} + +func checkResultPath(path string) error { + if !filepath.IsAbs(path) { + return errors.New("final result evidence path must be absolute") + } + info, err := os.Stat(filepath.Dir(path)) + if err != nil || !info.IsDir() || info.Mode().Perm() != 0700 { + return errors.New("final result evidence directory must be private mode 0700") + } + if _, err := os.Lstat(path); err == nil || !errors.Is(err, os.ErrNotExist) { + return errors.New("final result evidence already exists or is inaccessible") + } + return nil +} + +func persistOneShotResult(path string, body []byte) error { + if err := checkResultPath(path); err != nil { + return err + } + parent := filepath.Dir(path) + file, err := os.OpenFile(path, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0600) + if err != nil { + return errors.New("final result evidence already exists or cannot be created") + } + if _, err := file.Write(body); err != nil { + file.Close() + return err + } + if err := file.Sync(); err != nil { + file.Close() + return err + } + if err := file.Close(); err != nil { + return err + } + dir, err := os.Open(parent) + if err != nil { + return err + } + defer dir.Close() + return dir.Sync() +} + +type resultConsumer struct { + connection *amqp.Connection + channel *amqp.Channel + deliveries <-chan amqp.Delivery +} + +func openResultConsumer(brokerURL string) (_ *resultConsumer, err error) { + if err := requireDedicatedVhost(brokerURL); err != nil { + return nil, err + } + conn, err := amqp.Dial(brokerURL) + if err != nil { + return nil, fmt.Errorf("connect SaaS result queue: %T", err) + } + defer func() { + if err != nil { + conn.Close() + } + }() + ch, err := conn.Channel() + if err != nil { + return nil, fmt.Errorf("open SaaS result channel: %T", err) + } + queue, err := ch.QueueDeclarePassive(resultQueue, true, false, false, false, nil) + if err != nil { + ch.Close() + return nil, errors.New("preprovisioned SaaS result queue unavailable") + } + if queue.Messages != 0 || queue.Consumers != 0 { + ch.Close() + return nil, errors.New("SaaS result queue is not empty or has another consumer") + } + deliveries, err := ch.Consume(resultQueue, "", false, false, false, false, nil) + if err != nil { + ch.Close() + return nil, errors.New("SaaS result consumer unavailable") + } + return &resultConsumer{connection: conn, channel: ch, deliveries: deliveries}, nil +} + +func (c *resultConsumer) Close() { + if c != nil { + c.channel.Close() + c.connection.Close() + } +} + +func (c *resultConsumer) wait(ctx context.Context, path, dispatcherID, sourceID, taskID, callee, trunkID string) (oneShotResult, error) { + for { + select { + case <-ctx.Done(): + return oneShotResult{}, errors.New("approved call final result not received before timeout; outcome unknown") + case message, ok := <-c.deliveries: + if !ok { + return oneShotResult{}, errors.New("SaaS result subscription lost; outcome unknown") + } + summary, err := summarizeOneShotResult(message.Body, dispatcherID, sourceID, taskID, callee, trunkID) + if err != nil { + return oneShotResult{}, errors.Join(err, message.Nack(false, true)) + } + if err := persistOneShotResult(path, message.Body); err != nil { + return oneShotResult{}, errors.Join(err, message.Nack(false, true)) + } + if err := message.Ack(false); err != nil { + return oneShotResult{}, errors.New("final result persisted but SaaS MQ acknowledgment unknown") + } + return summary, nil + } + } +} + +func publishAndAwait(ctx context.Context, brokerURL string, data dataset, eventID, taskID, callee, resultPath string, at time.Time) (oneShotResult, error) { + if err := checkResultPath(resultPath); err != nil { + return oneShotResult{}, err + } + var task configread.Task + if err := json.Unmarshal(data.tasks[taskID], &task); err != nil || len(task.AllowedTrunkIDs) != 1 { + return oneShotResult{}, errors.New("one-shot result requires a single approved trunk") + } + consumer, err := openResultConsumer(brokerURL) + if err != nil { + return oneShotResult{}, err + } + defer consumer.Close() + publishCtx, cancel := context.WithTimeout(ctx, 10*time.Second) + defer cancel() + if err := publishExecuteAt(publishCtx, brokerURL, data, eventID, taskID, callee, at); err != nil { + return oneShotResult{}, err + } + return consumer.wait(ctx, resultPath, data.dispatcherID, eventID, taskID, callee, task.AllowedTrunkIDs[0]) +} diff --git a/scripts/check-current-mq-mock.sh b/scripts/check-current-mq-mock.sh index 9f56401..e972cf3 100644 --- a/scripts/check-current-mq-mock.sh +++ b/scripts/check-current-mq-mock.sh @@ -70,3 +70,4 @@ 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 +run_required_test ./deploys/test/saas-mock TestSaaSMockAwaitsExactResultBeforeStoppingCapture