From c61858fdedeb8e9e7437f42056b033fcaf34b54a Mon Sep 17 00:00:00 2001 From: Rogee Date: Tue, 29 Sep 2026 22:09:15 +0800 Subject: [PATCH] Apply approved dialogue limits in shared callflow --- .../saas-dispatcher-implementation.md | 5 +- internal/ai/current_limits_test.go | 83 +++++++++++++ internal/ai/current_pipeline.go | 59 ++++++++-- internal/callflow/approved.go | 40 +++++++ internal/callflow/approved_test.go | 109 ++++++++++++++++++ internal/callflow/capture.go | 12 +- internal/callflow/flow.go | 52 +++++---- internal/callflow/flow_test.go | 2 +- 8 files changed, 320 insertions(+), 42 deletions(-) create mode 100644 internal/ai/current_limits_test.go create mode 100644 internal/callflow/approved.go create mode 100644 internal/callflow/approved_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index ff96412..fa31922 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -51,8 +51,9 @@ - Dispatcher 在启动、增量发现、恢复任务、SIP 更新及发指令前对完整 AI/provider 快照执行能力校验;不支持的任务关闭准入而不占额度或呼叫。已按用户确认保留现行 TTS Schema:火山 TTS V2 SDK 无法表达的 PCMA、0.5–2 倍以外或整数 speech_rate 无法精确表达的速度均明确拒绝,不静默修改配置;不支持的插话配置同样拒绝。 - `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 不吻合即拒绝。 -- 验证:`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`go test -race ./internal/ai ./internal/configread ./internal/dispatcher ./internal/rpc -count=1`、`go vet ./...`、`go build ./...`、`git diff --check` 均通过。 -- **尚未完成:** 生产 CLI 接入、会话期间全部对话控制参数的实际媒体行为、完整录音/OSS/最终结果及真实 Agent/Asterisk、SaaS、MQ、AI 供应商联调。这些不能被本地 Mock 结果代签,分别属于 P05 收尾、P06–P08 或第二阶段外部验收。 +- 对话控制隔离验证:共享 `callflow` 不再根据转写或回复内的硬编码词推断拒联,仅处理明确的关键词/拒联事实;ASR-only 不播报开场或 TTS 回复。配置的首语音/静默期限、整通话期限与最大轮数交给媒体控制器;回复按 Unicode 字符数分片,火山 TTS 的缓存音频块数超限明确失败而不重试;不支持的插话配置在准入前拒绝。该控制器尚未接入新主 CLI,不能宣称真实通话媒体已验证。 +- 验证:`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`go test -race ./internal/ai ./internal/callflow ./internal/configread ./internal/dispatcher ./internal/rpc -count=1`、`go vet ./...`、`go build ./...`、`git diff --check` 均通过。 +- **尚未完成:** 新主 CLI 接入、对话控制与真实媒体的完整联动、录音/OSS/最终结果及真实 Agent/Asterisk、SaaS、MQ、AI 供应商联调。这些不能被本地 Mock 结果代签,分别属于 P05 收尾、P06–P08 或第二阶段外部验收。 ## 验收台账 diff --git a/internal/ai/current_limits_test.go b/internal/ai/current_limits_test.go new file mode 100644 index 0000000..e3ca19e --- /dev/null +++ b/internal/ai/current_limits_test.go @@ -0,0 +1,83 @@ +package ai + +import ( + "context" + "encoding/base64" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" +) + +func TestCurrentTTSRejectsExcessPendingAudioWithoutRetry(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + task = changeCurrentAgent(t, task, func(agent map[string]any) { + agent["conversation"].(map[string]any)["max_pending_audio_chunks"] = 1 + }) + var requests atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + requests.Add(1) + for _, pcm := range [][]byte{{1, 0}, {2, 0}} { + _, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString(pcm)) + } + _, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`) + })) + defer server.Close() + p := providers["tts-example"] + p.Endpoint = server.URL + providers[p.ProviderRef] = p + bound, err := BindCurrent(task, providers) + if err != nil { + t.Fatal(err) + } + if _, err := bound.Synthesize(context.Background(), "批准回复"); err == nil || !strings.Contains(err.Error(), "pending audio") || strings.Contains(err.Error(), p.Credential) { + t.Fatalf("explicit pending-audio bound was ignored or leaked credentials: %v", err) + } + if requests.Load() != 1 { + t.Fatalf("SDK must not automatically retry/bill twice: requests=%d", requests.Load()) + } +} + +func TestCurrentReplyChunksUnicodeByApprovedSentenceLimit(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + task = changeCurrentAgent(t, task, func(agent map[string]any) { + conversation := agent["conversation"].(map[string]any) + conversation["sentence_max_chars"] = 3 + conversation["max_pending_audio_chunks"] = 3 + }) + texts := make(chan string, 3) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + var payload struct { + Params struct { + Text string `json:"text"` + } `json:"req_params"` + } + if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + http.Error(w, "invalid SDK request", http.StatusBadRequest) + return + } + texts <- payload.Params.Text + _, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString([]byte{1, 0})) + _, _ = fmt.Fprintln(w, `{"code":20000000,"message":"ok","data":null}`) + })) + defer server.Close() + p := providers["tts-example"] + p.Endpoint = server.URL + providers[p.ProviderRef] = p + bound, err := BindCurrent(task, providers) + if err != nil { + t.Fatal(err) + } + audio, err := bound.SynthesizeReply(context.Background(), "你好世界,再见") + if err != nil || len(audio) != 6 { + t.Fatalf("reply must synthesize all bounded chunks: len=%d err=%v", len(audio), err) + } + for _, want := range []string{"你好世", "界,再", "见"} { + if got := <-texts; got != want { + t.Fatalf("Unicode reply was lost/reordered while splitting: got=%q want=%q", got, want) + } + } +} diff --git a/internal/ai/current_pipeline.go b/internal/ai/current_pipeline.go index 9d7ec6e..43eb90f 100644 --- a/internal/ai/current_pipeline.go +++ b/internal/ai/current_pipeline.go @@ -102,7 +102,7 @@ func (c *CurrentCall) RunTurn(ctx context.Context, pcm16 []byte) (TurnResult, er if err != nil { return TurnResult{}, fmt.Errorf("LLM failed: %w", err) } - audio, err := c.bound.Synthesize(ctx, reply) + audio, err := c.bound.SynthesizeReply(ctx, reply) if err != nil { return TurnResult{}, fmt.Errorf("TTS failed: %w", err) } @@ -186,12 +186,46 @@ func (b CurrentBound) Complete(ctx context.Context, finalUserText string) (strin } func (b CurrentBound) Synthesize(ctx context.Context, text string) ([]byte, error) { - if b.Mode != "full_ai" || b.TTS == nil { - return nil, errors.New("ASR-only mode does not call TTS") - } + audio, _, err := b.synthesize(ctx, text, 0) + return audio, err +} + +// SynthesizeReply splits a reply at the approved Unicode character limit, +// retaining every character and the same pending-audio bound across requests. +func (b CurrentBound) SynthesizeReply(ctx context.Context, text string) ([]byte, error) { if strings.TrimSpace(text) == "" { return nil, errors.New("TTS requires nonempty reply") } + maxChars := b.Conversation.SentenceMaxChars + if maxChars <= 0 { + return b.Synthesize(ctx, text) + } + runes := []rune(text) + var audio []byte + pending := 0 + for start := 0; start < len(runes); start += maxChars { + end := min(start+maxChars, len(runes)) + part, count, err := b.synthesize(ctx, string(runes[start:end]), pending) + if err != nil { + return nil, err + } + pending += count + audio = append(audio, part...) + } + return audio, nil +} + +func (b CurrentBound) synthesize(ctx context.Context, text string, alreadyPending int) ([]byte, int, error) { + if b.Mode != "full_ai" || b.TTS == nil { + return nil, 0, errors.New("ASR-only mode does not call TTS") + } + if strings.TrimSpace(text) == "" { + return nil, 0, errors.New("TTS requires nonempty reply") + } + limit := b.Conversation.MaxPendingAudioChunks + if limit > 0 && alreadyPending >= limit { + return nil, 0, errors.New("TTS exceeded approved pending audio chunk limit") + } ctx, cancel := currentDeadline(ctx, b.TTS.Timeout) defer cancel() request := b.TTS.Request // per-call copy: concurrent calls never share mutable SDK parameters @@ -201,24 +235,31 @@ func (b CurrentBound) Synthesize(ctx context.Context, text string) ([]byte, erro doubaospeech.WithBaseURL(b.TTS.Provider.Endpoint), ) var audio []byte + chunks := 0 completed := false for chunk, err := range client.TTSV2.Stream(ctx, &request) { if err != nil { - return nil, err + return nil, chunks, err } if chunk == nil { - return nil, errors.New("TTS returned an empty stream chunk") + return nil, chunks, errors.New("TTS returned an empty stream chunk") + } + if len(chunk.Audio) > 0 { + chunks++ + if limit > 0 && alreadyPending+chunks > limit { + return nil, chunks, errors.New("TTS exceeded approved pending audio chunk limit") + } + audio = append(audio, chunk.Audio...) } - audio = append(audio, chunk.Audio...) if chunk.IsLast { completed = true break } } if !completed || len(audio) == 0 || len(audio)%2 != 0 { - return nil, errors.New("TTS stream ended without complete PCM16 audio") + return nil, chunks, errors.New("TTS stream ended without complete PCM16 audio") } - return audio, nil + return audio, chunks, nil } func currentDeadline(ctx context.Context, timeout time.Duration) (context.Context, context.CancelFunc) { diff --git a/internal/callflow/approved.go b/internal/callflow/approved.go new file mode 100644 index 0000000..bd9f8c6 --- /dev/null +++ b/internal/callflow/approved.go @@ -0,0 +1,40 @@ +package callflow + +import ( + "context" + "errors" + + "git.ipao.vip/rogee/go-sip/internal/ai" +) + +// ApprovedPipeline is one call's immutable Agent execution. It shares the +// media sequencing with other modes without importing legacy AI snapshots. +type ApprovedPipeline interface { + Open(context.Context) ([]byte, error) + RunTurn(context.Context, []byte) (ai.TurnResult, error) +} + +func ExecuteApproved(ctx context.Context, session MediaSession, mode ai.Mode, pipeline ApprovedPipeline, capture CaptureConfig) (Result, error) { + if pipeline == nil { + return Result{}, errors.New("approved AI pipeline is required") + } + if capture.CallDuration > 0 { + var cancel context.CancelFunc + ctx, cancel = context.WithTimeout(ctx, capture.CallDuration) + defer cancel() + } + return executeFlow(ctx, session, mode, pipeline.Open, pipeline.RunTurn, capture) +} + +// ApprovedCapture hands explicit task conversation limits to the shared media +// controller. A zero field denotes an absent optional business value. +func ApprovedCapture(bound ai.CurrentBound) CaptureConfig { + return CaptureConfig{ + FirstSpeechTimeout: bound.Conversation.SilenceTimeout, + CallDuration: bound.Conversation.MaxDuration, + EndSilence: bound.Conversation.SilenceTimeout, + VoiceThreshold: 1, + MaxTurns: bound.Conversation.MaxTurns, + MaxPendingAudioChunks: bound.Conversation.MaxPendingAudioChunks, + } +} diff --git a/internal/callflow/approved_test.go b/internal/callflow/approved_test.go new file mode 100644 index 0000000..25b8ded --- /dev/null +++ b/internal/callflow/approved_test.go @@ -0,0 +1,109 @@ +package callflow + +import ( + "bytes" + "context" + "strings" + "testing" + "time" + + "git.ipao.vip/rogee/go-sip/internal/ai" +) + +type approvedFlowPipeline struct { + openingCalls int + turnCalls int + turn ai.TurnResult +} + +func (p *approvedFlowPipeline) Open(context.Context) ([]byte, error) { + p.openingCalls++ + return make([]byte, 6400), nil +} + +func (p *approvedFlowPipeline) RunTurn(context.Context, []byte) (ai.TurnResult, error) { + p.turnCalls++ + return p.turn, nil +} + +func TestApprovedKeywordStopsBeforeSendingReply(t *testing.T) { + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "我不用了", EndedByKeyword: true, 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.EndedByKeyword || len(result.Turns) != 1 || len(result.OutboundTurns) != 1 || session.stats.SentPackets != 1 || pipeline.openingCalls != 1 || pipeline.turnCalls != 1 { + t.Fatalf("keyword hangup must stop before LLM/TTS reply: turns=%d outbound=%d sent=%d", len(result.Turns), len(result.OutboundTurns), session.stats.SentPackets) + } +} + +func TestApprovedUnconfiguredRefusalTextDoesNotInventHangup(t *testing.T) { + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "不用了", Reply: "继续为您服务", 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: 1, + }) + if err != nil { + t.Fatal(err) + } + if result.Turn.InvalidCall || result.Turn.EndedByKeyword || len(result.OutboundTurns) != 2 || session.stats.SentPackets != 2 { + t.Fatalf("no unapproved hardcoded refusal match may terminate call: %+v", result.Turn) + } +} + +func TestApprovedASROnlySkipsOpeningAndReplyEvenIfPipelineHasAudio(t *testing.T) { + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "最终识别", AudioPCM16: make([]byte, 6400)}} + session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} + result, err := ExecuteApproved(context.Background(), session, ai.ModeASROnly, pipeline, CaptureConfig{ + FirstSpeechTimeout: time.Second, MaxDuration: 2 * time.Millisecond, MaxTurns: 1, + }) + if err != nil { + t.Fatal(err) + } + 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) + } +} + +func TestApprovedCapturePreservesConversationControls(t *testing.T) { + bound := ai.CurrentBound{Conversation: ai.CurrentConversation{SilenceTimeout: 3 * time.Second, MaxDuration: 2 * time.Minute, MaxTurns: 20, SentenceMaxChars: 80, MaxPendingAudioChunks: 32}} + capture := ApprovedCapture(bound) + if capture.FirstSpeechTimeout != 3*time.Second || capture.EndSilence != 3*time.Second || capture.CallDuration != 2*time.Minute || capture.MaxDuration != 0 || capture.MaxTurns != 20 || capture.VoiceThreshold <= 0 || capture.MaxPendingAudioChunks != 32 { + t.Fatalf("approved conversation limits were not handed to media controller: %+v", capture) + } +} + +type approvedDeadlineProbe struct{ deadlinePresent bool } + +func (*approvedDeadlineProbe) Open(context.Context) ([]byte, error) { return nil, nil } +func (p *approvedDeadlineProbe) RunTurn(ctx context.Context, _ []byte) (ai.TurnResult, error) { + deadline, ok := ctx.Deadline() + p.deadlinePresent = ok && time.Until(deadline) > 0 && time.Until(deadline) <= time.Second + return ai.TurnResult{Transcript: "最终识别", EndedByKeyword: true}, nil +} + +func TestApprovedConversationDurationBoundsWholeCall(t *testing.T) { + pipeline := &approvedDeadlineProbe{} + session := &scriptedTurnSession{turns: [][]byte{bytes.Repeat([]byte{1, 0}, 3200)}} + capture := ApprovedCapture(ai.CurrentBound{Conversation: ai.CurrentConversation{ + SilenceTimeout: 2 * time.Millisecond, MaxDuration: time.Second, MaxTurns: 1, + }}) + if _, err := ExecuteApproved(context.Background(), session, ai.ModeASROnly, pipeline, capture); err != nil || !pipeline.deadlinePresent { + t.Fatalf("approved overall duration was not enforced throughout media and AI: deadline=%t err=%v", pipeline.deadlinePresent, err) + } +} + +func TestApprovedSilenceTimeoutRejectsMissingSpeech(t *testing.T) { + pipeline := &approvedFlowPipeline{turn: ai.TurnResult{Transcript: "不应执行"}} + session := &scriptedTurnSession{turns: [][]byte{make([]byte, 6400)}} + capture := ApprovedCapture(ai.CurrentBound{Conversation: ai.CurrentConversation{ + SilenceTimeout: 10 * time.Millisecond, MaxDuration: 100 * time.Millisecond, MaxTurns: 1, + }}) + _, err := ExecuteApproved(context.Background(), session, ai.ModeASROnly, pipeline, capture) + if err == nil || !strings.Contains(err.Error(), "speech") || pipeline.turnCalls != 0 { + t.Fatalf("configured silence timeout did not stop before AI: turns=%d err=%v", pipeline.turnCalls, err) + } +} diff --git a/internal/callflow/capture.go b/internal/callflow/capture.go index 0cf2ea5..74013cf 100644 --- a/internal/callflow/capture.go +++ b/internal/callflow/capture.go @@ -11,11 +11,13 @@ import ( // collecting until bounded duration or end-of-speech silence, then send one // canonical PCM16 turn to the AI adapter. type CaptureConfig struct { - FirstSpeechTimeout time.Duration - MaxDuration time.Duration - EndSilence time.Duration - VoiceThreshold int - MaxTurns int + FirstSpeechTimeout time.Duration + MaxDuration time.Duration // maximum duration of one captured utterance + CallDuration time.Duration // maximum duration of the entire approved call + EndSilence time.Duration + VoiceThreshold int + MaxTurns int + MaxPendingAudioChunks int } func captureTurn(ctx context.Context, session MediaSession, cfg CaptureConfig) ([]byte, error) { diff --git a/internal/callflow/flow.go b/internal/callflow/flow.go index 29e6255..c77a206 100644 --- a/internal/callflow/flow.go +++ b/internal/callflow/flow.go @@ -4,7 +4,6 @@ import ( "context" "errors" "fmt" - "strings" "time" "git.ipao.vip/rogee/go-sip/internal/ai" @@ -44,22 +43,36 @@ func ExecuteWithCapture(ctx context.Context, session MediaSession, pipeline ai.P if session == nil || pipeline == nil { return Result{}, errors.New("media session and AI pipeline are required") } - if snapshot.Mode != ai.ModeFullAI && snapshot.Mode != ai.ModeASROnly { - return Result{}, fmt.Errorf("unsupported callflow AI mode %q", snapshot.Mode) + 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") + } + if mode != ai.ModeFullAI && mode != ai.ModeASROnly { + return Result{}, fmt.Errorf("unsupported callflow AI mode %q", mode) } maxTurns := capture.MaxTurns if maxTurns <= 0 { maxTurns = 1 } result := Result{} - if snapshot.Mode == ai.ModeFullAI { - openingPCM, err := pipeline.Synthesize(ctx, snapshot, opening) + if mode == ai.ModeFullAI { + openingPCM, err := open(ctx) if err != nil { return Result{}, err } - result.OutboundTurns = append(result.OutboundTurns, clonePCM(openingPCM)) - if err := session.SendPCM16(ctx, openingPCM, 16000); err != nil { - return result, err + if len(openingPCM) > 0 { + result.OutboundTurns = append(result.OutboundTurns, clonePCM(openingPCM)) + if err := session.SendPCM16(ctx, openingPCM, 16000); err != nil { + return result, err + } } } for turnIndex := 0; turnIndex < maxTurns; turnIndex++ { @@ -72,21 +85,23 @@ func ExecuteWithCapture(ctx context.Context, session MediaSession, pipeline ai.P if len(inbound) < 3200 { return result, fmt.Errorf("captured audio turn %d is too short", turnIndex+1) } - turn, err := pipeline.RunTurn(ctx, snapshot, inbound) + turn, err := runTurn(ctx, inbound) if err != nil { return result, fmt.Errorf("run AI turn %d: %w", turnIndex+1, err) } - turn.InvalidCall, turn.InvalidReason = invalidCallReason(turn.Transcript, turn.Reply) result.Turn = turn result.Turns = append(result.Turns, turn) - if turn.InvalidCall { + if turn.InvalidCall || turn.EndedByKeyword { result.RTP = session.Stats() return result, nil } - if snapshot.Mode == ai.ModeASROnly { + if mode == ai.ModeASROnly { result.RTP = session.Stats() continue } + if len(turn.AudioPCM16) == 0 { + return result, fmt.Errorf("AI reply turn %d has no audio", turnIndex+1) + } if err := session.SendPCM16(ctx, turn.AudioPCM16, 16000); err != nil { return result, fmt.Errorf("send AI reply turn %d: %w", turnIndex+1, err) } @@ -101,19 +116,6 @@ func clonePCM(pcm []byte) []byte { return append([]byte(nil), pcm...) } -func invalidCallReason(transcript, reply string) (bool, string) { - if strings.Contains(reply, ai.InvalidCallMarker) { - return true, "llm_invalid_call_marker" - } - normalized := strings.NewReplacer(" ", "", " ", "", "。", "", ",", "", ",", "", ".", "").Replace(strings.TrimSpace(transcript)) - for _, marker := range []string{"打错", "不需要", "不用", "没兴趣", "不考虑", "不方便", "别打", "拒绝", "骚扰", "语音信箱", "自动语音", "请按键", "空号"} { - if strings.Contains(normalized, marker) { - return true, "transcript_invalid_intent" - } - } - return false, "" -} - // MemorySession is a bounded transport adapter for mock/mixed-flow tests. It // has no SIP semantics and never bypasses the shared call flow. type MemorySession struct { diff --git a/internal/callflow/flow_test.go b/internal/callflow/flow_test.go index 62ba94f..9c7092d 100644 --- a/internal/callflow/flow_test.go +++ b/internal/callflow/flow_test.go @@ -90,7 +90,7 @@ func (invalidCallPipeline) Synthesize(context.Context, ai.Snapshot, string) ([]b } func (invalidCallPipeline) RunTurn(context.Context, ai.Snapshot, []byte) (ai.TurnResult, error) { - return ai.TurnResult{Transcript: "打错了", Reply: ai.InvalidCallMarker, AudioPCM16: make([]byte, 6400)}, nil + return ai.TurnResult{Transcript: "打错了", Reply: ai.InvalidCallMarker, AudioPCM16: make([]byte, 6400), InvalidCall: true, InvalidReason: "llm_invalid_call_marker"}, nil } func TestInvalidCallStopsBeforeReply(t *testing.T) {