diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 3337a09..2499f95 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -68,7 +68,7 @@ - 内存录音隔离组件:`RecordingSession` 仅复制共享通话流程实际读到和成功发送的 16-kHz PCM16,`EncodeMonoWAV` 直接在内存生成有界单声道 WAV;空音频、奇数字节、超过上限及未成功发送的音频都不能伪造成可上传录音。单元与 race 测试未产生业务文件。批准执行入口尚未接入该组件,且 Mock 中观测到的帧不等于真实 Asterisk 通话的全量媒体验收。 - Agent 录音交付隔离组件:`RecordingDelivery` 先确认结束,再依照录音是否实际生成分别上报唯一空录音结果或请求原授权并直传内存 WAV;录音生成失败保留通话真实结果、空录音对象及明确原因,不虚构上传事实。隔离测试通过本地 HTTP PUT 和假 Dispatcher RPC 覆盖成功无业务文件、OSS 明确失败后私有文件保存、恢复写入失败、未知 PUT 隔离、重启重领原目标、上传已确认后只重发原结果。再次调用不会隐式重新 PUT;正常已确认上传但尚未被 D 持久收讫的跨进程间隙仍受 K16 边界约束。此处未连接主入口批准执行媒体或 MQ;下述隔离链路另测本地真正的 D gRPC/SQLite。 - 最终结果隔离组件:`FinalResultPayload` 只从批准的任务身份、确认时间及实际收到的用户媒体窗口生成当前唯一结果;最终 ASR 文本与字面关键词拒联写入完整转写,未播放的助手候选内容不冒充转写,未知 SIP 响应码保留 `null`,无录音保留 `{}`。单元测试直接核验当前 MQ Schema、字段缺失/乱序时窗与无应答负例;开场白发送失败不得计入已播放音频。此组件尚未接入主入口,也未证明真实 RTP 转写时间精度。 -- 本地隔离链路:`RecordingSession` 从实际读出的 Mock 媒体生成内存 WAV,经 Agent→Dispatcher 双向 TLS gRPC 确认结束、领取原始资产的签名授权,再对本地 HTTPS OSS Mock 单次 PUT;Dispatcher 将原上传事实与唯一最终结果 outbox 同事务提交。测试确认一次 PUT、一次结果、无业务文件,使用的是官方 SDK 生成的路径,OSS Mock 不验证真实服务商签名;这不是主入口执行、RabbitMQ 投递或外部验收。 +- 本地隔离链路:共享批准通话流程 `ExecuteApproved` 让显式 ASR 测试桩按 ASR-only 限制消费实际读出的合成 PCM 媒体帧;`RecordingSession` 生成内存 WAV,`FinalResultPayload` 仅用最终模拟识别及实测采集时窗构造结果。Agent 经双向 TLS gRPC 向 Dispatcher 确认结束、领取原始资产签名授权,再对本地 HTTPS OSS Mock 单次 PUT;Dispatcher 将原上传事实与唯一最终结果 outbox 同事务提交。测试确认一次 PUT、一次结果、无业务文件,使用官方 SDK 生成的路径;ASR 测试桩及 OSS Mock 不验证真实供应商协议或签名。这不是主入口执行、RabbitMQ 投递或外部验收。 - 已验证:`go test ./... -count=1`、`go test -race ./... -count=1`、`go vet ./...`、`go build ./...`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`git diff --check`。尚未完成批准执行媒体录音到 Agent 交付组件的主入口接线、主入口 Agent↔Dispatcher 实际会话与录音执行接线(隔离 mTLS gRPC 已测)、真实执行时上传事实/最终结果交付及 MQ/端到端验收,不能宣称 P06 通过。 ## 验收台账 diff --git a/internal/rpc/recording_delivery_integration_test.go b/internal/rpc/recording_delivery_integration_test.go index 8fa598f..0abda71 100644 --- a/internal/rpc/recording_delivery_integration_test.go +++ b/internal/rpc/recording_delivery_integration_test.go @@ -4,6 +4,7 @@ import ( "bytes" "context" "encoding/json" + "errors" "io" "net" "net/http" @@ -16,6 +17,7 @@ import ( agentpb "git.ipao.vip/rogee/go-sip/gen/agent" "git.ipao.vip/rogee/go-sip/internal/agent" + "git.ipao.vip/rogee/go-sip/internal/ai" "git.ipao.vip/rogee/go-sip/internal/callflow" "git.ipao.vip/rogee/go-sip/internal/oss" @@ -24,15 +26,29 @@ import ( "google.golang.org/grpc/test/bufconn" ) +type isolatedApprovedASR struct{} + +func (isolatedApprovedASR) Open(context.Context) ([]byte, error) { return nil, nil } +func (isolatedApprovedASR) RunTurn(_ context.Context, pcm []byte) (ai.TurnResult, error) { + if len(pcm) != 3200 { + return ai.TurnResult{}, errors.New("isolated ASR did not receive the actual media frame") + } + return ai.TurnResult{Transcript: "隔离 Mock 最终识别文本"}, nil +} + func TestRecordingDeliveryRealMutualTLSAndLocalOSSCommitsOneSQLiteResult(t *testing.T) { server, database, _, request, snapshot := recordingRPCFixture(t) - mediaSession, err := callflow.NewRecordingSession(callflow.NewMemorySession(bytes.Repeat([]byte{1, 0}, 320)), 1024) + mediaSession, err := callflow.NewRecordingSession(callflow.NewMemorySession(bytes.Repeat([]byte{1, 0}, 1600)), 4096) if err != nil { t.Fatal(err) } - frame, err := mediaSession.ReadPayload(context.Background()) - if err != nil || len(frame) != 640 { - t.Fatalf("live media was not captured: size=%d err=%v", len(frame), err) + startedAt := time.Now().UTC() + callResult, err := callflow.ExecuteApproved(context.Background(), mediaSession, ai.ModeASROnly, isolatedApprovedASR{}, callflow.ApprovedCapture(ai.CurrentBound{ + Mode: "asr_only", Conversation: ai.CurrentConversation{SilenceTimeout: 100 * time.Millisecond, MaxDuration: time.Second, MaxTurns: 1}, + })) + endedAt := time.Now().UTC() + if err != nil || len(callResult.Turns) != 1 || len(callResult.OutboundTurns) != 0 || callResult.Turns[0].Transcript != "隔离 Mock 最终识别文本" { + t.Fatalf("approved ASR-only capture changed live media or generated unwanted audio: turns=%d outbound=%d err=%v", len(callResult.Turns), len(callResult.OutboundTurns), err) } wav, durationMS, err := mediaSession.WAV() if err != nil { @@ -41,7 +57,7 @@ func TestRecordingDeliveryRealMutualTLSAndLocalOSSCommitsOneSQLiteResult(t *test var puts atomic.Int32 localOSS := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { puts.Add(1) - body, err := io.ReadAll(io.LimitReader(r.Body, 1025)) + body, err := io.ReadAll(io.LimitReader(r.Body, 4097)) if err != nil || r.Method != http.MethodPut || !strings.HasPrefix(r.URL.Path, "/mock-bucket/approved/") || !bytes.Equal(body, wav) { t.Errorf("local OSS received an incorrect PUT: method=%q size=%d path=%q err=%v", r.Method, len(body), r.URL.Path, err) w.WriteHeader(http.StatusBadRequest) @@ -53,7 +69,7 @@ func TestRecordingDeliveryRealMutualTLSAndLocalOSSCommitsOneSQLiteResult(t *test server.OSS, err = oss.NewClient(oss.Config{ Endpoint: localOSS.URL, Region: "cn-test", Bucket: "mock-bucket", KeyPrefix: "approved", AccessKeyID: "isolated-test-key", AccessKeySecret: "isolated-test-secret", - GrantTTL: 15 * time.Minute, MaxAssetBytes: 1024, + GrantTTL: 15 * time.Minute, MaxAssetBytes: 4096, }) if err != nil { t.Fatal(err) @@ -92,17 +108,22 @@ func TestRecordingDeliveryRealMutualTLSAndLocalOSSCommitsOneSQLiteResult(t *test Session: func(context.Context) (*agentpb.RequestMeta, error) { return request.Meta, nil }, } - var result map[string]any - if err := json.Unmarshal(recordingResultPayload(t, snapshot, nil), &result); err != nil { - t.Fatal(err) - } - result["outcome"] = "answered" - result["reason_code"] = 200 - result["reason_message"] = "completed" - payload, err := json.Marshal(result) + payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{ + TaskID: snapshot.Task.TaskID, CallerProfileID: snapshot.Task.CallerProfileID, + Callee: "15003164745", TrunkID: "trunk-mock", StartedAt: startedAt, EndedAt: endedAt, + Outcome: "answered", ReasonMessage: "isolated Mock media completed", + }, callResult) if err != nil { t.Fatal(err) } + var result struct { + Transcript []struct { + Text string `json:"text"` + } `json:"transcript"` + } + if err := json.Unmarshal(payload, &result); err != nil || len(result.Transcript) != 1 || result.Transcript[0].Text != callResult.Turns[0].Transcript { + t.Fatalf("final result lost actual final ASR text: segments=%d err=%v", len(result.Transcript), err) + } root := t.TempDir() if err := os.Chmod(root, 0700); err != nil { t.Fatal(err)