diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 97fe981..4647e01 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -52,7 +52,8 @@ - `ExecuteApproved` 只在隔离 Mock 中签发:携带任务原始 JSON、引用的 provider 明文凭据、SIP revision、选定路由/主叫/原始被叫、独立的拨号期限及最大通话时限;任务+provider 原始字节+SIP revision 计算绑定摘要。Agent 对照已激活的 D 会话身份、绑定摘要、SDK 能力及**实际观察到的** SIP 加载结果;发出 Mock 指令前仅持久保存摘要和未知占用,失败/重复/重启不自动重拨,文件中不留明文凭据。当前 `LoadedSIP` 为可注入的隔离 Mock 观察器,**不等于真实 Asterisk 已加载核验**。 - ASR-only 不启动 LLM/TTS;full-AI 使用冻结的 model、显式 temperature=0/max_tokens、TTS voice/format/speech_rate 与 provider 原值凭据。真实 ASR SDK WebSocket 发出的音频输入/识别配置、LLM 与 TTS SDK 向隔离 HTTP Mock 发出的实际参数均已捕获核对;已配置的开场白由 TTS 单次合成,失败不自动重播且未完成时不进入对话轮次;空开场白不调用 TTS。关键词仅匹配用户侧最终 ASR,失败/未知挂断不自动重试。`ApprovedOriginator` 经生成的 Unary gRPC Stub 交付完整快照,SIP 全量 revision 不吻合即拒绝。 - 对话控制隔离验证:共享 `callflow` 不再根据转写或回复内的硬编码词推断拒联,仅处理明确的关键词/拒联事实;ASR-only 不播报开场或 TTS 回复。配置的首语音/静默期限、整通话期限与最大轮数交给媒体控制器;回复按 Unicode 字符数分片,火山 TTS 的缓存音频块数超限明确失败而不重试;不支持的插话配置在准入前拒绝。该控制器尚未接入新主 CLI,不能宣称真实通话媒体已验证。 -- Agent 隔离媒体入口 `rpc.RunApprovedCall` 仅接收签发的执行快照、媒体会话、挂断动作和显式每通话 AI pipeline;缺失或 typed nil 适配器在读媒体前拒绝,无隐式 SDK/Mock 回退。测试分别由已绑定快照构建 SDK pipeline 与注入隔离 ASR 测试桩。签发通话期限与 AI 总期限均约束整通电话;ASR-only 不合成开场,full-AI 开场完成才采集语音。父级期限不会被“未检测到语音”掩盖,空媒体 Mock 等到上下文真正结束。当前媒体只处理 16 kHz PCM16,SDK 虽接受、但媒体不能正确处理的 24 kHz ASR/TTS 配置在 Dispatcher 与 Agent 绑定时明确拒绝,不静默播放或识别错速音频。该入口尚未连接主 CLI,也未用于真实拨号。 +- Agent 隔离媒体入口 `rpc.RunApprovedCall` 仅接收签发的执行快照、媒体会话、挂断动作和显式每通话 AI pipeline;缺失或 typed nil 适配器在读媒体前拒绝,无隐式 SDK/Mock 回退。测试分别由已绑定快照构建 SDK pipeline 与注入受控合成脚本的 `ApprovedMockPipeline`。签发通话期限与 AI 总期限均约束整通电话;ASR-only 不合成开场,full-AI 开场完成才采集语音。父级期限不会被“未检测到语音”掩盖,空媒体 Mock 等到上下文真正结束。当前媒体只处理 16 kHz PCM16,SDK 虽接受、但媒体不能正确处理的 24 kHz ASR/TTS 配置在 Dispatcher 与 Agent 绑定时明确拒绝,不静默播放或识别错速音频。该入口尚未连接主 CLI,也未用于真实拨号。 +- 隔离 Mock AI 适配器仅消费显式冻结的模拟媒体与最终识别脚本,不手写或替代供应商协议;缺失/多余轮次、未批准的开场或助理音频明确拒绝。ASR-only 不执行开场与回复;full-AI 开场只允许一次,用户最终识别的已配置字面关键词先挂断、绝不播放该轮候选回复,未配置关键词不推断拒联。任务未配置最大轮数时严格使用共享通话流程的一轮边界,负数仍拒绝。Mock 输出不代表真实 ASR/LLM/TTS 已通过。 - 验证:`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`、`go test ./... -count=1`、`go test -race ./internal/ai ./internal/callflow ./internal/configread ./internal/dispatcher ./internal/rpc -count=1`、`go vet ./...`、`go build ./...`、`git diff --check` 均通过。 - **后续边界:** 录音直传/OSS 失败恢复及最终结果属 P06;新主 CLI 接线、真实媒体完整联动及旧路径清理属 P07;隔离端到端和 A01–A12 属 P08。真实 Agent/Asterisk、SaaS、MQ、AI 供应商联调未开展,不能由 Mock 结果代签。 @@ -68,7 +69,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。 - 最终结果隔离组件:`ApprovedExecution` 已保留获批任务的 `caller_profile_id`;`FinalResultPayload` 只从批准的任务身份、确认时间及实际收到的用户媒体窗口生成当前唯一结果;最终 ASR 文本与字面关键词拒联写入完整转写,未播放的助手候选内容不冒充转写,未知 SIP 响应码保留 `null`,无录音保留 `{}`。单元测试直接核验当前 MQ Schema、字段缺失/乱序时窗与无应答负例;开场白发送失败不得计入已播放音频。此组件尚未接入主入口,也未证明真实 RTP 转写时间精度。 -- 本地隔离链路:共享批准通话流程 `ExecuteApproved` 让显式 ASR 测试桩按 ASR-only 限制消费实际读出的合成 PCM 媒体帧;`RecordingSession` 生成内存 WAV,`FinalResultPayload` 仅用最终模拟识别及实测采集时窗构造结果。Agent 经双向 TLS gRPC 向 Dispatcher 确认结束、领取原始资产签名授权,再对本地 HTTPS OSS Mock 单次 PUT;Dispatcher 将原上传事实与唯一最终结果 outbox 同事务提交。测试确认一次 PUT、一次结果、无业务文件,使用官方 SDK 生成的路径;ASR 测试桩及 OSS Mock 不验证真实供应商协议或签名。这不是主入口执行、RabbitMQ 投递或外部验收。 +- 本地隔离链路:`RunApprovedCall` 显式使用合成 `ApprovedMockPipeline`,让共享批准通话流程按 ASR-only 限制消费实际读出的 PCM Mock 媒体帧;`RecordingSession` 生成内存 WAV,`FinalResultPayload` 仅用最终模拟识别及实测采集时窗构造结果。Agent 经双向 TLS gRPC 向 Dispatcher 确认结束、领取原始资产签名授权,再对本地 HTTPS OSS Mock 单次 PUT;Dispatcher 将原上传事实与唯一最终结果 outbox 同事务提交。测试确认一次 PUT、一次结果、无业务文件,使用官方 SDK 生成的路径;Mock ASR 与 OSS 服务不验证真实供应商协议或签名。这不是主入口执行、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/ai/approved_mock.go b/internal/ai/approved_mock.go new file mode 100644 index 0000000..d6f2cb4 --- /dev/null +++ b/internal/ai/approved_mock.go @@ -0,0 +1,142 @@ +package ai + +import ( + "bytes" + "context" + "errors" + "strings" + "sync" + "unicode/utf8" +) + +// ApprovedMockScript is explicit synthetic media and final ASR text for local +// isolation only. It is never a substitute for an actual provider response. +type ApprovedMockScript struct { + OpeningPCM16 []byte + Turns []ApprovedMockTurn +} + +type ApprovedMockTurn struct { + Transcript string + Reply string + ReplyPCM16 []byte +} + +// ApprovedMockPipeline is private to one approved Mock call. The signed AI +// parameters remain owned by CurrentBound; this adapter neither contacts a +// provider nor silently manufactures absent script turns or audio. +type ApprovedMockPipeline struct { + mu sync.Mutex + mode string + openingPCM []byte + turns []ApprovedMockTurn + keyword *KeywordHangup + opened bool + ended bool + next int +} + +func NewApprovedMockPipeline(bound CurrentBound, script ApprovedMockScript, hangup func(context.Context) error) (*ApprovedMockPipeline, error) { + keyword, err := NewKeywordHangup(bound.HangupKeywords, hangup) + if err != nil { + return nil, err + } + maxTurns := bound.Conversation.MaxTurns + if maxTurns == 0 { + maxTurns = 1 // executeFlow uses one turn when this optional limit is absent. + } + if maxTurns < 1 || len(script.Turns) != maxTurns { + return nil, errors.New("Mock script must cover exactly the approved maximum turns") + } + switch bound.Mode { + case string(ModeASROnly): + if bound.Opening != "" || bound.LLM != nil || bound.TTS != nil || len(script.OpeningPCM16) != 0 { + return nil, errors.New("ASR-only Mock cannot synthesize opening or assistant audio") + } + case string(ModeFullAI): + if bound.LLM == nil || bound.TTS == nil { + return nil, errors.New("full-AI Mock requires approved LLM and TTS configuration") + } + if (bound.Opening == "" && len(script.OpeningPCM16) != 0) || (bound.Opening != "" && !validMockPCM(script.OpeningPCM16)) { + return nil, errors.New("Mock opening audio must match the approved opening") + } + default: + return nil, errors.New("Mock AI mode is unsupported") + } + frozen := make([]ApprovedMockTurn, len(script.Turns)) + for i, turn := range script.Turns { + if !utf8.ValidString(turn.Transcript) || !utf8.ValidString(turn.Reply) { + return nil, errors.New("Mock text is not UTF-8") + } + if bound.Mode == string(ModeASROnly) { + if turn.Reply != "" || len(turn.ReplyPCM16) != 0 { + return nil, errors.New("ASR-only Mock contains assistant audio") + } + } else { + keywordTurn := false + for _, literal := range bound.HangupKeywords { + if strings.Contains(turn.Transcript, literal) { + keywordTurn = true + break + } + } + if !keywordTurn && (turn.Reply == "" || !validMockPCM(turn.ReplyPCM16)) { + return nil, errors.New("full-AI Mock reply is missing or invalid") + } + if len(turn.ReplyPCM16)%2 != 0 { + return nil, errors.New("Mock reply PCM16 has an odd byte count") + } + } + frozen[i] = ApprovedMockTurn{Transcript: turn.Transcript, Reply: turn.Reply, ReplyPCM16: bytes.Clone(turn.ReplyPCM16)} + } + return &ApprovedMockPipeline{mode: bound.Mode, openingPCM: bytes.Clone(script.OpeningPCM16), turns: frozen, keyword: keyword}, nil +} + +func validMockPCM(pcm []byte) bool { return len(pcm) > 0 && len(pcm)%2 == 0 } + +func (p *ApprovedMockPipeline) Open(ctx context.Context) ([]byte, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + p.mu.Lock() + defer p.mu.Unlock() + if p.opened { + return nil, errors.New("Mock opening already attempted") + } + p.opened = true + if p.mode == string(ModeASROnly) { + return nil, nil + } + return bytes.Clone(p.openingPCM), nil +} + +func (p *ApprovedMockPipeline) RunTurn(ctx context.Context, pcm []byte) (TurnResult, error) { + if err := ctx.Err(); err != nil { + return TurnResult{}, err + } + p.mu.Lock() + defer p.mu.Unlock() + // The shared ASR-only flow never calls Open; full-AI must play its opening first. + if (p.mode == string(ModeFullAI) && !p.opened) || p.ended || p.next >= len(p.turns) { + return TurnResult{}, errors.New("Mock turn was not approved or has already ended") + } + if len(pcm) < 3200 || len(pcm)%2 != 0 { + return TurnResult{}, errors.New("Mock ASR requires observed 16-kHz PCM16 media") + } + turn := p.turns[p.next] + p.next++ // A transport or hangup failure cannot replay this turn. + stopped, err := p.keyword.Handle(ctx, ASRSegment{Source: "user", Text: turn.Transcript, Final: true}) + if stopped { + p.ended = true + } + if err != nil { + return TurnResult{}, err + } + if stopped { + return TurnResult{Transcript: turn.Transcript, EndedByKeyword: true}, nil + } + if p.mode == string(ModeASROnly) { + return TurnResult{Transcript: turn.Transcript}, nil + } + return TurnResult{Transcript: turn.Transcript, Reply: turn.Reply, AudioPCM16: bytes.Clone(turn.ReplyPCM16)}, nil +} diff --git a/internal/ai/approved_mock_test.go b/internal/ai/approved_mock_test.go new file mode 100644 index 0000000..c91289e --- /dev/null +++ b/internal/ai/approved_mock_test.go @@ -0,0 +1,93 @@ +package ai + +import ( + "bytes" + "context" + "testing" +) + +func TestApprovedMockASROnlyDoesNotInventRefusalOrAssistantAudio(t *testing.T) { + var hangups int + pipeline, err := NewApprovedMockPipeline(CurrentBound{Mode: "asr_only", Conversation: CurrentConversation{MaxTurns: 1}}, ApprovedMockScript{ + Turns: []ApprovedMockTurn{{Transcript: "我拒绝但是没有配置关键词"}}, + }, func(context.Context) error { hangups++; return nil }) + if err != nil { + t.Fatal(err) + } + if opening, err := pipeline.Open(context.Background()); err != nil || len(opening) != 0 { + t.Fatalf("ASR-only mock generated an opening: size=%d err=%v", len(opening), err) + } + turn, err := pipeline.RunTurn(context.Background(), bytes.Repeat([]byte{1, 0}, 1600)) + if err != nil || turn.Transcript != "我拒绝但是没有配置关键词" || turn.Reply != "" || len(turn.AudioPCM16) != 0 || turn.EndedByKeyword || hangups != 0 { + t.Fatalf("mock invented unapproved full-AI or refusal facts: turn=%+v hangups=%d err=%v", turn, hangups, err) + } + if _, err := pipeline.RunTurn(context.Background(), []byte{1, 0}); err == nil { + t.Fatal("mock replayed an already used ASR result") + } +} + +func TestApprovedMockAbsentOptionalTurnLimitUsesSharedOneTurnBoundary(t *testing.T) { + pipeline, err := NewApprovedMockPipeline(CurrentBound{Mode: "asr_only"}, ApprovedMockScript{Turns: []ApprovedMockTurn{{Transcript: "隔离识别"}}}, func(context.Context) error { return nil }) + if err != nil { + t.Fatal(err) + } + if turn, err := pipeline.RunTurn(context.Background(), bytes.Repeat([]byte{1, 0}, 1600)); err != nil || turn.Transcript != "隔离识别" { + t.Fatalf("absent optional turn limit incorrectly rejected one approved turn: turn=%+v err=%v", turn, err) + } +} + +func TestApprovedMockFullAIFreezesScriptAndStopsBeforeKeywordReply(t *testing.T) { + opening := bytes.Repeat([]byte{1, 0}, 320) + reply := bytes.Repeat([]byte{2, 0}, 320) + script := ApprovedMockScript{OpeningPCM16: opening, Turns: []ApprovedMockTurn{ + {Transcript: "正常模拟用户语音", Reply: "模拟助手答复", ReplyPCM16: reply}, + {Transcript: "请不要联系", Reply: "不得播报", ReplyPCM16: reply}, + }} + var hangups int + bound := CurrentBound{Mode: "full_ai", Opening: "批准的开场", HangupKeywords: []string{"不要联系"}, LLM: &CurrentLLM{}, TTS: &CurrentTTS{}, Conversation: CurrentConversation{MaxTurns: 2}} + pipeline, err := NewApprovedMockPipeline(bound, script, func(context.Context) error { hangups++; return nil }) + if err != nil { + t.Fatal(err) + } + opening[0], reply[0], script.Turns[0].Transcript = 9, 9, "修改后的转写" + played, err := pipeline.Open(context.Background()) + if err != nil || len(played) != 640 || played[0] != 1 { + t.Fatalf("opening was not frozen from the approved mock fixture: size=%d err=%v", len(played), err) + } + if _, err := pipeline.Open(context.Background()); err == nil { + t.Fatal("ambiguous opening was replayed") + } + first, err := pipeline.RunTurn(context.Background(), bytes.Repeat([]byte{1, 0}, 1600)) + if err != nil || first.Transcript != "正常模拟用户语音" || first.Reply != "模拟助手答复" || len(first.AudioPCM16) != 640 || first.AudioPCM16[0] != 2 { + t.Fatalf("full-AI mock changed approved fixture bytes: turn=%+v err=%v", first, err) + } + last, err := pipeline.RunTurn(context.Background(), bytes.Repeat([]byte{1, 0}, 1600)) + if err != nil || !last.EndedByKeyword || last.Transcript != "请不要联系" || last.Reply != "" || len(last.AudioPCM16) != 0 || hangups != 1 { + t.Fatalf("keyword did not stop before the assistant reply: turn=%+v hangups=%d err=%v", last, hangups, err) + } +} + +func TestApprovedMockRejectsMissingAIOrUnscriptedMedia(t *testing.T) { + for _, tc := range []struct { + name string + bound CurrentBound + script ApprovedMockScript + }{ + {"missing_script", CurrentBound{Mode: "asr_only", Conversation: CurrentConversation{MaxTurns: 1}}, ApprovedMockScript{}}, + {"asr_with_llm", CurrentBound{Mode: "asr_only", LLM: &CurrentLLM{}, Conversation: CurrentConversation{MaxTurns: 1}}, ApprovedMockScript{Turns: []ApprovedMockTurn{{Transcript: "测试"}}}}, + {"full_without_llm", CurrentBound{Mode: "full_ai", TTS: &CurrentTTS{}, Conversation: CurrentConversation{MaxTurns: 1}}, ApprovedMockScript{Turns: []ApprovedMockTurn{{Transcript: "测试"}}}}, + {"unavailable_opening_audio", CurrentBound{Mode: "full_ai", LLM: &CurrentLLM{}, TTS: &CurrentTTS{}, Opening: "批准开场", Conversation: CurrentConversation{MaxTurns: 1}}, ApprovedMockScript{Turns: []ApprovedMockTurn{{Transcript: "测试", Reply: "模拟答复", ReplyPCM16: []byte{1, 0}}}}}, + {"unsolicited_odd_opening_audio", CurrentBound{Mode: "full_ai", LLM: &CurrentLLM{}, TTS: &CurrentTTS{}, Conversation: CurrentConversation{MaxTurns: 1}}, ApprovedMockScript{OpeningPCM16: []byte{1}, Turns: []ApprovedMockTurn{{Transcript: "测试", Reply: "模拟答复", ReplyPCM16: []byte{1, 0}}}}}, + {"extra_unapproved_turn", CurrentBound{Mode: "asr_only", Conversation: CurrentConversation{MaxTurns: 1}}, ApprovedMockScript{Turns: []ApprovedMockTurn{{Transcript: "第一轮"}, {Transcript: "第二轮"}}}}, + {"negative_turn_limit", CurrentBound{Mode: "asr_only", Conversation: CurrentConversation{MaxTurns: -1}}, ApprovedMockScript{Turns: []ApprovedMockTurn{{Transcript: "测试"}}}}, + } { + t.Run(tc.name, func(t *testing.T) { + if _, err := NewApprovedMockPipeline(tc.bound, tc.script, func(context.Context) error { return nil }); err == nil { + t.Fatal("unapproved Mock settings were accepted") + } + }) + } + if _, err := NewApprovedMockPipeline(CurrentBound{Mode: "asr_only"}, ApprovedMockScript{Turns: []ApprovedMockTurn{{Transcript: "测试"}}}, nil); err == nil { + t.Fatal("mock without a real hangup callback was admitted") + } +} diff --git a/internal/rpc/approved_runner_test.go b/internal/rpc/approved_runner_test.go index 0a3331c..839f4b8 100644 --- a/internal/rpc/approved_runner_test.go +++ b/internal/rpc/approved_runner_test.go @@ -119,11 +119,14 @@ func TestApprovedCallRunnerUsesExplicitIsolatedPipelineAndObservedMedia(t *testi if err != nil { t.Fatal(err) } - approved := ApprovedExecution{ - AI: ai.CurrentBound{Mode: "asr_only", Conversation: ai.CurrentConversation{SilenceTimeout: 100 * time.Millisecond, MaxDuration: time.Second, MaxTurns: 1}}, - MaxCallDuration: time.Second, + bound := ai.CurrentBound{Mode: "asr_only", Conversation: ai.CurrentConversation{SilenceTimeout: 100 * time.Millisecond, MaxDuration: time.Second, MaxTurns: 1}} + hangup := func(context.Context) error { return nil } + mock, err := ai.NewApprovedMockPipeline(bound, ai.ApprovedMockScript{Turns: []ai.ApprovedMockTurn{{Transcript: "隔离 Mock 最终识别文本"}}}, hangup) + if err != nil { + t.Fatal(err) } - result, err := RunApprovedCall(context.Background(), approved, mediaSession, func(context.Context) error { return nil }, isolatedApprovedASR{}) + approved := ApprovedExecution{AI: bound, MaxCallDuration: time.Second} + result, err := RunApprovedCall(context.Background(), approved, mediaSession, hangup, mock) if err != nil || len(result.Turns) != 1 || len(result.CaptureWindows) != 1 || len(result.OutboundTurns) != 0 || result.Turns[0].Transcript != "隔离 Mock 最终识别文本" { t.Fatalf("explicit isolated pipeline did not honor approved ASR-only media: turns=%d windows=%d outbound=%d err=%v", len(result.Turns), len(result.CaptureWindows), len(result.OutboundTurns), err) } @@ -140,7 +143,7 @@ func TestApprovedCallRunnerRejectsTypedNilPipelineBeforeMedia(t *testing.T) { t.Fatalf("missing pipeline was allowed to consume call media: err=%v stats=%+v", err, session.Stats()) } var missingMedia *callflow.MemorySession - _, err = RunApprovedCall(context.Background(), ApprovedExecution{AI: ai.CurrentBound{Mode: "asr_only"}, MaxCallDuration: time.Second}, missingMedia, func(context.Context) error { return nil }, isolatedApprovedASR{}) + _, err = RunApprovedCall(context.Background(), ApprovedExecution{AI: ai.CurrentBound{Mode: "asr_only"}, MaxCallDuration: time.Second}, missingMedia, func(context.Context) error { return nil }, &ai.CurrentCall{}) if err == nil { t.Fatal("typed-nil media session was admitted") } diff --git a/internal/rpc/recording_delivery_integration_test.go b/internal/rpc/recording_delivery_integration_test.go index 0abda71..c4ec1e7 100644 --- a/internal/rpc/recording_delivery_integration_test.go +++ b/internal/rpc/recording_delivery_integration_test.go @@ -4,7 +4,6 @@ import ( "bytes" "context" "encoding/json" - "errors" "io" "net" "net/http" @@ -26,26 +25,22 @@ 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}, 1600)), 4096) if err != nil { t.Fatal(err) } + bound := ai.CurrentBound{Mode: "asr_only", Conversation: ai.CurrentConversation{SilenceTimeout: 100 * time.Millisecond, MaxDuration: time.Second}} + hangup := func(context.Context) error { return nil } + mock, err := ai.NewApprovedMockPipeline(bound, ai.ApprovedMockScript{Turns: []ai.ApprovedMockTurn{{Transcript: "隔离 Mock 最终识别文本"}}}, hangup) + if err != nil { + t.Fatal(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}, - })) + callResult, err := RunApprovedCall(context.Background(), ApprovedExecution{ + AI: bound, MaxCallDuration: time.Second, TaskID: snapshot.Task.TaskID, CallerProfileID: snapshot.Task.CallerProfileID, + }, mediaSession, hangup, mock) 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)