refactor(callflow): retire unapproved AI execution entry

This commit is contained in:
2026-09-30 14:32:38 +08:00
parent 1eb1e8d93b
commit 4ef20d69c4
4 changed files with 94 additions and 232 deletions
@@ -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 实际采集区间,再迁移共用测试夹具;现行获批通话测试和全仓测试通过。未声称真实媒体链路通过。
## 验收台账
+91
View File
@@ -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)
}
+2 -26
View File
@@ -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")
-206
View File
@@ -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)
}
}