From 4ef20d69c4de5936c192bcfb32d274796ddde072 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 14:32:38 +0800 Subject: [PATCH] refactor(callflow): retire unapproved AI execution entry --- .../saas-dispatcher-implementation.md | 1 + internal/callflow/approved_test.go | 91 ++++++++ internal/callflow/flow.go | 28 +-- internal/callflow/flow_test.go | 206 ------------------ 4 files changed, 94 insertions(+), 232 deletions(-) delete mode 100644 internal/callflow/flow_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index ffe7739..50b1990 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -124,6 +124,7 @@ - Agent Proto 旧消息:先用描述符结构测试复现 `ConfigReference` 等废弃定义仍可被发现,再按当前八个服务方法及现行 Go 引用追踪字段依赖;删除无调用者的 26 个旧消息和 3 个旧枚举,不改现行请求/响应的字段号。使用本地 Buf 重新生成 Go 类型,核对七项来源/hash;现行录音客户端、获批执行、控制及会话测试保持通过。固定 `agent.v1` 仍是有效的内部协议值;这不构成真实 Agent/Asterisk 或外部 SaaS 验收。 - 隔离 MQ 校验脚本:审计发现脚本因测试文件改名仍使用旧筛选式,`internal/mq` 和 `internal/dispatcher` 输出 `no tests to run` 但退出成功;先复现两个空匹配,再改为现行测试名并要求三项均明确 `PASS`,跳过或空匹配均失败。隔离 RabbitMQ 实跑 `TestBrokerSharedResultQueueAndNoConfigure`、`TestRuntimeIsolatedControlBacklogExecuteAndSharedResult` 与 `TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent` 全部通过;这只证明本机隔离链路,不代替真实 MQ/SaaS 应用收讫。 - 旧 ARI/媒体执行包:`internal/callruntime` 仅实现旧 `ai.Snapshot` 的独立 `Run` 入口,无当前 Agent/Dispatcher 调用者;删除该包及仅针对该路径的测试。现行获批执行、录音恢复和完整 AI Mock 测试继续覆盖唯一现行入口;删除死代码不表示真实 Asterisk/ARI 或线路已验证,历史验收报告保留当时包名与覆盖率事实。 +- 旧通话流程入口:旧 `Execute`/`ExecuteWithCapture` 使用可缺省的旧 AI 快照,现已无 Agent 调用者;删除入口及专属测试,保留同一 `executeFlow` 媒体顺序实现。先在当前 `ExecuteApproved` 的测试中补齐首次媒体等待上限、三轮对话、开场发送失败不伪造播放、拒联后不发送回答与 ASR-only 实际采集区间,再迁移共用测试夹具;现行获批通话测试和全仓测试通过。未声称真实媒体链路通过。 ## 验收台账 diff --git a/internal/callflow/approved_test.go b/internal/callflow/approved_test.go index 25b8ded..784bd20 100644 --- a/internal/callflow/approved_test.go +++ b/internal/callflow/approved_test.go @@ -8,6 +8,7 @@ import ( "time" "git.ipao.vip/rogee/go-sip/internal/ai" + "git.ipao.vip/rogee/go-sip/internal/media" ) type approvedFlowPipeline struct { @@ -26,6 +27,91 @@ func (p *approvedFlowPipeline) RunTurn(context.Context, []byte) (ai.TurnResult, return p.turn, nil } +type scriptedTurnSession struct { + turns [][]byte + index int + served bool + stats media.RTPStats +} + +func (s *scriptedTurnSession) ReadPayload(ctx context.Context) ([]byte, error) { + if !s.served && s.index < len(s.turns) { + s.served = true + payload := append([]byte(nil), s.turns[s.index]...) + s.index++ + s.stats.ReceivedPackets++ + s.stats.ReceivedBytes += uint64(len(payload)) + return payload, nil + } + <-ctx.Done() + s.served = false + return nil, ctx.Err() +} + +func (s *scriptedTurnSession) SendPCM16(ctx context.Context, pcm []byte, _ int) error { + if err := ctx.Err(); err != nil { + return err + } + s.stats.SentPackets++ + s.stats.SentBytes += uint64(len(pcm)) + return nil +} + +func (s *scriptedTurnSession) Stats() media.RTPStats { return s.stats } + +type rejectedOpeningSession struct{ MediaSession } + +func (rejectedOpeningSession) SendPCM16(context.Context, []byte, int) error { + return context.Canceled +} + +func TestApprovedInitialMediaWaitIsBounded(t *testing.T) { + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "不应执行"}} + _, err := ExecuteApproved(context.Background(), NewMemorySession(nil), ai.ModeASROnly, pipeline, CaptureConfig{ + FirstSpeechTimeout: time.Millisecond, MaxDuration: time.Millisecond, MaxTurns: 1, + }) + if err == nil || pipeline.turnCalls != 0 { + t.Fatalf("missing speech was not bounded: calls=%d err=%v", pipeline.turnCalls, err) + } +} + +func TestApprovedSharedFlowRunsThreeConversationTurns(t *testing.T) { + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "最终识别", Reply: "获批回答", AudioPCM16: make([]byte, 6400)}} + session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400), make([]byte, 6400), make([]byte, 6400)}} + result, err := ExecuteApproved(context.Background(), session, ai.ModeFullAI, pipeline, CaptureConfig{ + FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 3, + }) + if err != nil { + t.Fatal(err) + } + if pipeline.openingCalls != 1 || pipeline.turnCalls != 3 || len(result.Turns) != 3 || len(result.InboundTurns) != 3 || len(result.OutboundTurns) != 4 || result.RTP.ReceivedBytes != 19200 { + t.Fatalf("approved three-turn flow incomplete: opening=%d turns=%d inbound=%d outbound=%d received=%d", pipeline.openingCalls, pipeline.turnCalls, len(result.InboundTurns), len(result.OutboundTurns), result.RTP.ReceivedBytes) + } +} + +func TestApprovedOpeningSendFailureDoesNotInventPlayback(t *testing.T) { + session := rejectedOpeningSession{MediaSession: NewMemorySession(nil)} + pipeline := &approvedFlowPipeline{} + result, err := ExecuteApproved(context.Background(), session, ai.ModeFullAI, pipeline, CaptureConfig{MaxTurns: 1}) + if err == nil || len(result.OutboundTurns) != 0 || result.RTP.SentPackets != 0 || pipeline.turnCalls != 0 { + t.Fatalf("failed opening was reported as played: err=%v outbound=%d sent=%d turns=%d", err, len(result.OutboundTurns), result.RTP.SentPackets, pipeline.turnCalls) + } +} + +func TestApprovedInvalidCallStopsBeforeReply(t *testing.T) { + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "打错了", InvalidCall: true, InvalidReason: "llm_invalid_call_marker", AudioPCM16: make([]byte, 6400)}} + session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} + result, err := ExecuteApproved(context.Background(), session, ai.ModeFullAI, pipeline, CaptureConfig{ + FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 3, + }) + if err != nil { + t.Fatal(err) + } + if !result.Turn.InvalidCall || result.Turn.InvalidReason != "llm_invalid_call_marker" || len(result.Turns) != 1 || session.stats.SentPackets != 1 { + t.Fatalf("invalid call did not stop before reply: turn=%+v turns=%d sent=%d", result.Turn, len(result.Turns), session.stats.SentPackets) + } +} + func TestApprovedKeywordStopsBeforeSendingReply(t *testing.T) { pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "我不用了", EndedByKeyword: true, AudioPCM16: make([]byte, 6400)}} session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} @@ -57,12 +143,17 @@ func TestApprovedUnconfiguredRefusalTextDoesNotInventHangup(t *testing.T) { func TestApprovedASROnlySkipsOpeningAndReplyEvenIfPipelineHasAudio(t *testing.T) { pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "最终识别", AudioPCM16: make([]byte, 6400)}} session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} + before := time.Now() result, err := ExecuteApproved(context.Background(), session, ai.ModeASROnly, pipeline, CaptureConfig{ FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 1, }) + after := time.Now() if err != nil { t.Fatal(err) } + if len(result.CaptureWindows) != 1 || result.CaptureWindows[0].StartedAt.Before(before) || result.CaptureWindows[0].EndedAt.After(after) || result.CaptureWindows[0].EndedAt.Before(result.CaptureWindows[0].StartedAt) { + t.Fatalf("ASR-only has no observed media interval: %+v", result.CaptureWindows) + } if pipeline.openingCalls != 0 || pipeline.turnCalls != 1 || len(result.OutboundTurns) != 0 || session.stats.SentPackets != 0 { t.Fatalf("ASR-only cannot play opening or TTS reply: opening=%d turns=%d sent=%d", pipeline.openingCalls, pipeline.turnCalls, session.stats.SentPackets) } diff --git a/internal/callflow/flow.go b/internal/callflow/flow.go index c00a14c..a937664 100644 --- a/internal/callflow/flow.go +++ b/internal/callflow/flow.go @@ -10,9 +10,8 @@ import ( "git.ipao.vip/rogee/go-sip/internal/media" ) -// MediaSession is the only transport boundary of the shared call flow. -// Real mode uses Asterisk ExternalMedia/RTP; mock mode uses MemorySession; the -// sequencing is shared while ASR-only deliberately omits opening/reply TTS. +// MediaSession is the transport boundary of the approved call flow. +// The current isolated Mock uses MemorySession; ASR-only omits opening/reply TTS. type MediaSession interface { ReadPayload(context.Context) ([]byte, error) SendPCM16(context.Context, []byte, int) error @@ -35,29 +34,6 @@ type Result struct { RTP media.RTPStats } -func Execute(ctx context.Context, session MediaSession, pipeline ai.Pipeline, snapshot ai.Snapshot, opening string, turnWindow time.Duration) (Result, error) { - if turnWindow <= 0 { - turnWindow = 5 * time.Second - } - return ExecuteWithCapture(ctx, session, pipeline, snapshot, opening, CaptureConfig{ - FirstSpeechTimeout: turnWindow, - MaxDuration: turnWindow, - MaxTurns: 1, - }) -} - -func ExecuteWithCapture(ctx context.Context, session MediaSession, pipeline ai.Pipeline, snapshot ai.Snapshot, opening string, capture CaptureConfig) (Result, error) { - if session == nil || pipeline == nil { - return Result{}, errors.New("media session and AI pipeline are required") - } - return executeFlow(ctx, session, snapshot.Mode, - func(ctx context.Context) ([]byte, error) { return pipeline.Synthesize(ctx, snapshot, opening) }, - func(ctx context.Context, pcm []byte) (ai.TurnResult, error) { - return pipeline.RunTurn(ctx, snapshot, pcm) - }, - capture) -} - func executeFlow(ctx context.Context, session MediaSession, mode ai.Mode, open func(context.Context) ([]byte, error), runTurn func(context.Context, []byte) (ai.TurnResult, error), capture CaptureConfig) (Result, error) { if session == nil || open == nil || runTurn == nil { return Result{}, errors.New("media session and AI pipeline are required") diff --git a/internal/callflow/flow_test.go b/internal/callflow/flow_test.go deleted file mode 100644 index df75101..0000000 --- a/internal/callflow/flow_test.go +++ /dev/null @@ -1,206 +0,0 @@ -package callflow - -import ( - "context" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/contracts" - "git.ipao.vip/rogee/go-sip/internal/ai" - "git.ipao.vip/rogee/go-sip/internal/media" -) - -func TestInitialMediaWaitIsBounded(t *testing.T) { - raw, err := contracts.Read("examples/agent-version-full-explicit.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := ai.ValidateForMode(raw, ai.ModeFullAI) - if err != nil { - t.Fatal(err) - } - _, err = Execute(context.Background(), NewMemorySession(nil), ai.MockPipeline{}, snapshot, "测试开场", time.Millisecond) - if err == nil { - t.Fatal("expected bounded initial media wait to fail") - } -} - -func TestSharedFlowRunsThreeConversationTurns(t *testing.T) { - raw, err := contracts.Read("examples/agent-version-full-explicit.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := ai.ValidateForMode(raw, ai.ModeFullAI) - if err != nil { - t.Fatal(err) - } - session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400), make([]byte, 6400), make([]byte, 6400)}} - result, err := ExecuteWithCapture(context.Background(), session, ai.MockPipeline{MaxAudioBytes: 16 << 20}, snapshot, "测试开场", CaptureConfig{ - FirstSpeechTimeout: time.Second, - MaxDuration: 2 * time.Millisecond, - MaxTurns: 3, - }) - if err != nil { - t.Fatal(err) - } - if len(result.Turns) != 3 || len(result.InboundTurns) != 3 || len(result.OutboundTurns) != 4 { - t.Fatalf("expected opening plus three shared turns, got turns=%d inbound=%d outbound=%d", len(result.Turns), len(result.InboundTurns), len(result.OutboundTurns)) - } - if result.Turn.Transcript == "" || result.Turn.Reply == "" || result.RTP.ReceivedBytes != 19200 { - t.Fatalf("three-turn flow facts are incomplete: %+v", result) - } -} - -type scriptedTurnSession struct { - turns [][]byte - index int - served bool - stats media.RTPStats -} - -func (s *scriptedTurnSession) ReadPayload(ctx context.Context) ([]byte, error) { - if !s.served && s.index < len(s.turns) { - s.served = true - payload := append([]byte(nil), s.turns[s.index]...) - s.index++ - s.stats.ReceivedPackets++ - s.stats.ReceivedBytes += uint64(len(payload)) - return payload, nil - } - <-ctx.Done() - s.served = false - return nil, ctx.Err() -} - -func (s *scriptedTurnSession) SendPCM16(ctx context.Context, pcm []byte, _ int) error { - if err := ctx.Err(); err != nil { - return err - } - s.stats.SentPackets++ - s.stats.SentBytes += uint64(len(pcm)) - return nil -} - -func (s *scriptedTurnSession) Stats() media.RTPStats { return s.stats } - -type invalidCallPipeline struct{} - -func (invalidCallPipeline) Synthesize(context.Context, ai.Snapshot, string) ([]byte, error) { - return make([]byte, 6400), nil -} - -func (invalidCallPipeline) RunTurn(context.Context, ai.Snapshot, []byte) (ai.TurnResult, error) { - return ai.TurnResult{Transcript: "打错了", Reply: ai.InvalidCallMarker, AudioPCM16: make([]byte, 6400), InvalidCall: true, InvalidReason: "llm_invalid_call_marker"}, nil -} - -type rejectedOpeningSession struct{ MediaSession } - -func (rejectedOpeningSession) SendPCM16(context.Context, []byte, int) error { - return context.Canceled -} - -func TestFailedOpeningDoesNotInventOutboundAudio(t *testing.T) { - session := rejectedOpeningSession{MediaSession: NewMemorySession(nil)} - result, err := ExecuteWithCapture(context.Background(), session, invalidCallPipeline{}, ai.Snapshot{Mode: ai.ModeFullAI}, "approved opening", CaptureConfig{MaxTurns: 1}) - if err == nil || len(result.OutboundTurns) != 0 || result.RTP.SentPackets != 0 { - t.Fatalf("failed media send was reported as played audio: err=%v outbound=%d rtp=%+v", err, len(result.OutboundTurns), result.RTP) - } -} - -func TestInvalidCallStopsBeforeReply(t *testing.T) { - raw, err := contracts.Read("examples/agent-version-full-explicit.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := ai.ValidateForMode(raw, ai.ModeFullAI) - if err != nil { - t.Fatal(err) - } - session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} - result, err := ExecuteWithCapture(context.Background(), session, invalidCallPipeline{}, snapshot, "开场", CaptureConfig{ - FirstSpeechTimeout: time.Second, - MaxDuration: 2 * time.Millisecond, - MaxTurns: 3, - }) - if err != nil { - t.Fatal(err) - } - if !result.Turn.InvalidCall || result.Turn.InvalidReason != "llm_invalid_call_marker" || len(result.Turns) != 1 { - t.Fatalf("invalid call was not classified: %+v", result) - } - if session.stats.SentPackets != 1 { - t.Fatalf("invalid call must not send a reply, sent packets=%d", session.stats.SentPackets) - } -} - -func TestASROnlySkipsOpeningLLMAndTTS(t *testing.T) { - raw, err := contracts.Read("examples/agent-version-asr-only.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := ai.ValidateForMode(raw, ai.ModeASROnly) - if err != nil { - t.Fatal(err) - } - pipeline := &asrOnlyPipeline{} - session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} - before := time.Now() - result, err := ExecuteWithCapture(context.Background(), session, pipeline, snapshot, "不得播放", CaptureConfig{ - FirstSpeechTimeout: time.Second, - MaxDuration: time.Millisecond, - MaxTurns: 1, - }) - after := time.Now() - if err != nil { - t.Fatal(err) - } - if len(result.CaptureWindows) != 1 || result.CaptureWindows[0].StartedAt.Before(before) || result.CaptureWindows[0].EndedAt.After(after) || result.CaptureWindows[0].EndedAt.Before(result.CaptureWindows[0].StartedAt) { - t.Fatalf("final user ASR has no actually observed media interval: %+v", result.CaptureWindows) - } - if pipeline.synthesizeCalls != 0 || pipeline.turnCalls != 1 { - t.Fatalf("ASR-only provider calls = synthesize:%d turn:%d", pipeline.synthesizeCalls, pipeline.turnCalls) - } - if result.Turn.Transcript != "asr-only transcript" || result.Turn.Reply != "" || len(result.OutboundTurns) != 0 { - t.Fatalf("unexpected ASR-only result: %+v", result) - } - if session.stats.SentPackets != 0 { - t.Fatalf("ASR-only flow must not send opening/reply audio, sent packets=%d", session.stats.SentPackets) - } -} - -type asrOnlyPipeline struct { - synthesizeCalls int - turnCalls int -} - -func (p *asrOnlyPipeline) Synthesize(context.Context, ai.Snapshot, string) ([]byte, error) { - p.synthesizeCalls++ - return nil, nil -} - -func (p *asrOnlyPipeline) RunTurn(context.Context, ai.Snapshot, []byte) (ai.TurnResult, error) { - p.turnCalls++ - return ai.TurnResult{Transcript: "asr-only transcript"}, nil -} - -func TestMockModeUsesTheSharedConversationFlow(t *testing.T) { - raw, err := contracts.Read("examples/agent-version-full-explicit.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := ai.ValidateForMode(raw, ai.ModeFullAI) - if err != nil { - t.Fatal(err) - } - session := NewMemorySession(make([]byte, 6400)) - result, err := Execute(context.Background(), session, ai.MockPipeline{MaxAudioBytes: 16 << 20}, snapshot, "测试开场", 10*time.Millisecond) - if err != nil { - t.Fatal(err) - } - if result.Turn.Transcript == "" || result.Turn.Reply == "" || len(result.Turn.AudioPCM16) == 0 { - t.Fatalf("shared flow did not complete mock ASR/LLM/TTS: %+v", result.Turn) - } - if result.RTP.ReceivedBytes != 6400 || result.RTP.SentPackets == 0 || len(session.OutboundPCM()) == 0 { - t.Fatalf("shared media boundary was not exercised: %+v", result.RTP) - } -}