feat(test): await and persist one confirmed SaaS call result

This commit is contained in:
2026-10-04 17:11:57 +08:00
parent df51a825b0
commit 99a8c77132
5 changed files with 359 additions and 3 deletions
+1 -1
View File
@@ -13,7 +13,7 @@
- 准备测试 HTTPS 证书与私钥;将 `SAAS_MOCK_DISPATCHER_SECRET` 和含凭据的 `SAAS_MOCK_RABBITMQ_URL` 放在受限环境文件,不在命令行、仓库或聊天中传输。
- 启动:`go run ./deploys/test/saas-mock --data <私有目录> --dispatcher-id <UUID> --listen <地址:端口> --tls-cert <证书文件> --tls-key <私钥文件>`。
- 服务以标准 `X-DISPATCHER-id` 和 `X-DISPATCHER-SECRET-KEY` 校验 Dispatcher,再提供五类只读配置。错误的归属、资源、租户、快照或消息队列会导致拒绝启动/读取。
- 单次投递另起命令:`go run ./deploys/test/saas-mock --data <私有目录> --dispatcher-id <UUID> --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 <UUID> --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。
+16 -2
View File
@@ -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 {
+151
View File
@@ -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 == "" {
+190
View File
@@ -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])
}
+1
View File
@@ -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