Apply approved dialogue limits in shared callflow
This commit is contained in:
@@ -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 或第二阶段外部验收。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
+27
-25
@@ -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 {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user