From c45a05756d8d7850dfeadb62c8f6eeec80f8a6a1 Mon Sep 17 00:00:00 2001 From: Rogee Date: Sun, 4 Oct 2026 17:57:04 +0800 Subject: [PATCH] fix: distinguish SaaS dispatch receipt from final call result --- deploys/test/saas-mock/README.md | 4 +-- deploys/test/saas-mock/publish_test.go | 41 ++++++++++++++++++++++++++ deploys/test/saas-mock/result_wait.go | 34 +++++++++++++++++++-- 3 files changed, 74 insertions(+), 5 deletions(-) diff --git a/deploys/test/saas-mock/README.md b/deploys/test/saas-mock/README.md index b0d0b66..869028a 100644 --- a/deploys/test/saas-mock/README.md +++ b/deploys/test/saas-mock/README.md @@ -13,8 +13,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 <原始白名单号码> --await-result-file <私有结果文件>`。此命令在专用结果队列中等待精确匹配的单通最终结果,先将原始结果写入 `0600` 私有文件并同步磁盘,才确认 MQ 消费;标准输出只显示结果摘要/hash,不输出转写、录音或签名 URL。必须在 `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。`call.execute` 的 `dispatched` 只是派发回执:归属匹配时先私密保存为 `<结果文件>.receipt.json` 并确认,再继续等待唯一最终结果并保持抓包;不匹配的消息不确认。发布确认只代表 MQ 接收,不代表 SaaS 已收到最终结果;结果等待超时/归属不符时不清理未知通话,也不重发同通命令。已有未交付队列消息须人工确认处置,不自动清理。 测试:`go test ./deploys/test/saas-mock` 验证正式配置客户端;设置指向**单独隔离 vhost** 的 `SAAS_MOCK_TEST_BROKER_URL` 后,`TestSaaSMockProvisionsDispatcherTopology` 还将实际预建 MQ 并用 Dispatcher 被动读回。缺省测试不会连接共享 RabbitMQ。 -正式入口已经具备独立 `nonprod-real` 代码路径,但测试机的正式链路、真实 AI/TTS、录音、OSS 与拨号前抓包仍未完成整体验证。六组外呼均未执行。此服务通过测试也不代表真实 SaaS 或生产签收。 +正式入口已经具备独立 `nonprod-real` 代码路径,但测试机的正式链路、真实 AI/TTS、录音、OSS 与拨号前抓包仍未完成整体验证。2026-10-04 两次已确认的数企试拨操作均无 SIP 包;第二次产生派发回执但 Agent 失败,原因未查明,也未取得最终结果。此服务通过测试不代表真实 SaaS 或生产签收。 diff --git a/deploys/test/saas-mock/publish_test.go b/deploys/test/saas-mock/publish_test.go index 0f46eb5..f8b0caf 100644 --- a/deploys/test/saas-mock/publish_test.go +++ b/deploys/test/saas-mock/publish_test.go @@ -16,6 +16,36 @@ import ( amqp "github.com/rabbitmq/amqp091-go" ) +func TestOneShotDispatchReceiptDoesNotEndCapture(t *testing.T) { + data, err := loadDataset(testDataDir(t), testDispatcher) + if err != nil { + t.Fatal(err) + } + receipt := map[string]any{ + "event_id": "event-once-1", "event_type": "call.execute", + "dispatcher_id": data.dispatcherID, "tenant_id": data.tenantID, + "issued_at": "2026-10-04T02:00:00Z", "payload": map[string]any{"status": "dispatched"}, + } + body, err := json.Marshal(receipt) + if err != nil { + t.Fatal(err) + } + if err := contract.ValidateCurrent("mq", body); err != nil { + t.Fatal(err) + } + if err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "event-once-1"); err != nil { + t.Fatal(err) + } + if err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "different-event"); err == nil { + t.Fatal("unrelated call receipt must remain unacknowledged") + } + receipt["payload"].(map[string]any)["status"] = "unknown" + body, _ = json.Marshal(receipt) + if err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "event-once-1"); err == nil { + t.Fatal("unknown receipt status cannot authorize completion") + } +} + func TestOneShotResultMustMatchSourceIdentityAndPersistBeforeReceipt(t *testing.T) { data, err := loadDataset(testDataDir(t), testDispatcher) if err != nil { @@ -100,6 +130,10 @@ func TestSaaSMockAwaitsExactResultBeforeStoppingCapture(t *testing.T) { if err != nil { t.Fatal(err) } + receipt, err := json.Marshal(map[string]any{"event_id": eventID, "event_type": "call.execute", "dispatcher_id": data.dispatcherID, "tenant_id": data.tenantID, "issued_at": "2026-10-04T02:00:00Z", "payload": map[string]any{"status": "dispatched"}}) + if err != nil { + t.Fatal(err) + } ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() producer := make(chan error, 1) @@ -132,6 +166,10 @@ func TestSaaSMockAwaitsExactResultBeforeStoppingCapture(t *testing.T) { producer <- err return } + if err := ch.PublishWithContext(ctx, "agent-call.saas.v1", "d."+testDispatcher+".out", true, false, amqp.Publishing{DeliveryMode: amqp.Persistent, ContentType: "application/json", Body: receipt}); 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 @@ -162,6 +200,9 @@ func TestSaaSMockAwaitsExactResultBeforeStoppingCapture(t *testing.T) { if saved, err := os.ReadFile(file); err != nil || !bytes.Equal(saved, result) { t.Fatalf("result was not saved before broker receipt: err=%v", err) } + if saved, err := os.ReadFile(file + ".receipt.json"); err != nil || !bytes.Equal(saved, receipt) { + t.Fatalf("issuance receipt was not persisted before acknowledgment: err=%v", err) + } } func TestSaaSMockPublishesOneCommandOnProvisionedRabbitMQ(t *testing.T) { diff --git a/deploys/test/saas-mock/result_wait.go b/deploys/test/saas-mock/result_wait.go index 0bb4ebe..f553b44 100644 --- a/deploys/test/saas-mock/result_wait.go +++ b/deploys/test/saas-mock/result_wait.go @@ -29,6 +29,25 @@ func expectedResultEventID(dispatcherID, sourceID string) string { return fmt.Sprintf("result-%x", identity[:]) } +func verifyOneShotReceipt(body []byte, dispatcherID string, tenantID int64, sourceID string) error { + if err := contract.ValidateCurrent("mq", body); err != nil { + return errors.New("call dispatch receipt violates the current MQ contract") + } + var event struct { + EventID string `json:"event_id"` + EventType string `json:"event_type"` + DispatcherID string `json:"dispatcher_id"` + TenantID int64 `json:"tenant_id"` + Payload struct { + Status string `json:"status"` + } `json:"payload"` + } + if err := json.Unmarshal(body, &event); err != nil || event.EventID != sourceID || event.EventType != "call.execute" || event.DispatcherID != dispatcherID || event.TenantID != tenantID || event.Payload.Status != "dispatched" { + return errors.New("dispatch receipt does not belong to the selected call") + } + return nil +} + 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") @@ -144,7 +163,7 @@ func (c *resultConsumer) Close() { } } -func (c *resultConsumer) wait(ctx context.Context, path, dispatcherID, sourceID, taskID, callee, trunkID string) (oneShotResult, error) { +func (c *resultConsumer) wait(ctx context.Context, path, dispatcherID string, tenantID int64, sourceID, taskID, callee, trunkID string) (oneShotResult, error) { for { select { case <-ctx.Done(): @@ -155,7 +174,16 @@ func (c *resultConsumer) wait(ctx context.Context, path, dispatcherID, sourceID, } summary, err := summarizeOneShotResult(message.Body, dispatcherID, sourceID, taskID, callee, trunkID) if err != nil { - return oneShotResult{}, errors.Join(err, message.Nack(false, true)) + if verifyOneShotReceipt(message.Body, dispatcherID, tenantID, sourceID) != nil { + return oneShotResult{}, errors.Join(err, message.Nack(false, true)) + } + if err := persistOneShotResult(path+".receipt.json", message.Body); err != nil { + return oneShotResult{}, errors.Join(err, message.Nack(false, true)) + } + if err := message.Ack(false); err != nil { + return oneShotResult{}, errors.New("dispatch receipt persisted but SaaS MQ acknowledgment unknown") + } + continue // issuance receipt is not the unique physical call result } if err := persistOneShotResult(path, message.Body); err != nil { return oneShotResult{}, errors.Join(err, message.Nack(false, true)) @@ -186,5 +214,5 @@ func publishAndAwait(ctx context.Context, brokerURL string, data dataset, eventI 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]) + return consumer.wait(ctx, resultPath, data.dispatcherID, data.tenantID, eventID, taskID, callee, task.AllowedTrunkIDs[0]) }