acknowledge explicit no-dial rejections without waiting for a call result
This commit is contained in:
@@ -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 <原始白名单号码> --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 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` 并确认,再继续等待唯一最终结果并保持抓包。明确未拨号的 `rejected` 回执则先私密保存为 `<结果文件>.rejected.json` 再确认,并立即按无呼叫失败结束等待;不伪造最终通话结果。不匹配的消息不确认。发布确认只代表 MQ 接收,不代表 SaaS 已收到最终结果;结果等待超时/归属不符时不清理未知通话,也不重发同通命令。已有未交付队列消息须人工确认处置,不自动清理。
|
||||
|
||||
测试:`go test ./deploys/test/saas-mock` 验证正式配置客户端;设置指向**单独隔离 vhost** 的 `SAAS_MOCK_TEST_BROKER_URL` 后,`TestSaaSMockProvisionsDispatcherTopology` 还将实际预建 MQ 并用 Dispatcher 被动读回。缺省测试不会连接共享 RabbitMQ。
|
||||
|
||||
|
||||
@@ -33,19 +33,66 @@ func TestOneShotDispatchReceiptDoesNotEndCapture(t *testing.T) {
|
||||
if err := contract.ValidateCurrent("mq", body); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "event-once-1"); err != nil {
|
||||
if status, err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "event-once-1"); err != nil || status != "dispatched" {
|
||||
t.Fatalf("expected dispatch receipt: %s %v", status, 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"] = "rejected"
|
||||
receipt["payload"].(map[string]any)["reason_code"] = nil
|
||||
receipt["payload"].(map[string]any)["reason_message"] = "Agent refused before outbound call"
|
||||
body, _ = json.Marshal(receipt)
|
||||
if err := contract.ValidateCurrent("mq", body); 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")
|
||||
if status, err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "event-once-1"); err != nil || status != "rejected" {
|
||||
t.Fatalf("definite no-dial rejection should be acknowledged without a final result: %s %v", status, err)
|
||||
}
|
||||
receipt["payload"].(map[string]any)["status"] = "unknown"
|
||||
body, _ = json.Marshal(receipt)
|
||||
if err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "event-once-1"); err == nil {
|
||||
if _, err := verifyOneShotReceipt(body, data.dispatcherID, data.tenantID, "event-once-1"); err == nil {
|
||||
t.Fatal("unknown receipt status cannot authorize completion")
|
||||
}
|
||||
}
|
||||
|
||||
type testMQAck struct{ accepted, requeued int }
|
||||
|
||||
func (a *testMQAck) Ack(uint64, bool) error { a.accepted++; return nil }
|
||||
func (a *testMQAck) Nack(uint64, bool, bool) error { a.requeued++; return nil }
|
||||
func (a *testMQAck) Reject(uint64, bool) error { a.requeued++; return nil }
|
||||
|
||||
func TestRejectedOneShotReceiptIsSavedBeforeAckAndStopsWaiting(t *testing.T) {
|
||||
data, err := loadDataset(testDataDir(t), testDispatcher)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
body, err := json.Marshal(map[string]any{"event_id": "call-rejected", "event_type": "call.execute", "dispatcher_id": data.dispatcherID, "tenant_id": data.tenantID, "issued_at": "2026-10-04T02:00:00Z", "payload": map[string]any{"status": "rejected", "reason_code": nil, "reason_message": "Agent refused before outbound call"}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ack := &testMQAck{}
|
||||
incoming := make(chan amqp.Delivery, 1)
|
||||
incoming <- amqp.Delivery{Body: body, Acknowledger: ack}
|
||||
consumer := &resultConsumer{deliveries: incoming}
|
||||
dir := t.TempDir()
|
||||
if err := os.Chmod(dir, 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
file := filepath.Join(dir, "result.json")
|
||||
_, err = consumer.wait(context.Background(), file, data.dispatcherID, data.tenantID, "call-rejected", "task-full", "15003164745", "trunk-mock")
|
||||
if err == nil || !strings.Contains(err.Error(), "rejected") || ack.accepted != 1 || ack.requeued != 0 {
|
||||
t.Fatalf("rejected receipt was not acknowledged as terminal no-dial: accepted=%d requeued=%d err=%v", ack.accepted, ack.requeued, err)
|
||||
}
|
||||
saved, err := os.ReadFile(file + ".rejected.json")
|
||||
if err != nil || !bytes.Equal(saved, body) {
|
||||
t.Fatalf("receipt not durably retained: %v", err)
|
||||
}
|
||||
if _, err := os.Stat(file); !os.IsNotExist(err) {
|
||||
t.Fatalf("false final result created: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOneShotResultMustMatchSourceIdentityAndPersistBeforeReceipt(t *testing.T) {
|
||||
data, err := loadDataset(testDataDir(t), testDispatcher)
|
||||
if err != nil {
|
||||
|
||||
@@ -29,9 +29,9 @@ func expectedResultEventID(dispatcherID, sourceID string) string {
|
||||
return fmt.Sprintf("result-%x", identity[:])
|
||||
}
|
||||
|
||||
func verifyOneShotReceipt(body []byte, dispatcherID string, tenantID int64, sourceID string) error {
|
||||
func verifyOneShotReceipt(body []byte, dispatcherID string, tenantID int64, sourceID string) (string, error) {
|
||||
if err := contract.ValidateCurrent("mq", body); err != nil {
|
||||
return errors.New("call dispatch receipt violates the current MQ contract")
|
||||
return "", errors.New("call dispatch receipt violates the current MQ contract")
|
||||
}
|
||||
var event struct {
|
||||
EventID string `json:"event_id"`
|
||||
@@ -42,10 +42,10 @@ func verifyOneShotReceipt(body []byte, dispatcherID string, tenantID int64, sour
|
||||
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")
|
||||
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" && event.Payload.Status != "rejected") {
|
||||
return "", errors.New("dispatch receipt does not belong to the selected call")
|
||||
}
|
||||
return nil
|
||||
return event.Payload.Status, nil
|
||||
}
|
||||
|
||||
func summarizeOneShotResult(body []byte, dispatcherID, sourceID, taskID, callee, trunkID string) (oneShotResult, error) {
|
||||
@@ -174,15 +174,23 @@ func (c *resultConsumer) wait(ctx context.Context, path, dispatcherID string, te
|
||||
}
|
||||
summary, err := summarizeOneShotResult(message.Body, dispatcherID, sourceID, taskID, callee, trunkID)
|
||||
if err != nil {
|
||||
if verifyOneShotReceipt(message.Body, dispatcherID, tenantID, sourceID) != nil {
|
||||
status, receiptErr := verifyOneShotReceipt(message.Body, dispatcherID, tenantID, sourceID)
|
||||
if receiptErr != nil {
|
||||
return oneShotResult{}, errors.Join(err, message.Nack(false, true))
|
||||
}
|
||||
if err := persistOneShotResult(path+".receipt.json", message.Body); err != nil {
|
||||
receiptPath := path + ".receipt.json"
|
||||
if status == "rejected" {
|
||||
receiptPath = path + ".rejected.json"
|
||||
}
|
||||
if err := persistOneShotResult(receiptPath, 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")
|
||||
}
|
||||
if status == "rejected" {
|
||||
return oneShotResult{}, errors.New("Agent explicitly rejected this call before dialing; no final call result")
|
||||
}
|
||||
continue // issuance receipt is not the unique physical call result
|
||||
}
|
||||
if err := persistOneShotResult(path, message.Body); err != nil {
|
||||
|
||||
Reference in New Issue
Block a user