fix: distinguish SaaS dispatch receipt from final call result

This commit is contained in:
2026-10-04 17:57:04 +08:00
parent 99a8c77132
commit c45a05756d
3 changed files with 74 additions and 5 deletions
+2 -2
View File
@@ -13,8 +13,8 @@
- 准备测试 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 <原始白名单号码> --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 <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。`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 或生产签收。
+41
View File
@@ -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) {
+31 -3
View File
@@ -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])
}