From d9ad4204804ca9e48c0f11adf50e8082224ee501 Mon Sep 17 00:00:00 2001 From: Rogee Date: Tue, 29 Sep 2026 21:37:41 +0800 Subject: [PATCH] Require one-shot opening before full-AI turns --- .../saas-dispatcher-implementation.md | 2 +- internal/ai/current_opening_test.go | 138 ++++++++++++++++++ internal/ai/current_pipeline.go | 45 +++++- .../rpc/approved_full_ai_integration_test.go | 16 +- 4 files changed, 194 insertions(+), 7 deletions(-) create mode 100644 internal/ai/current_opening_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 0d9c5b6..ff96412 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -50,7 +50,7 @@ - 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 发出的实际参数均已捕获核对;关键词仅匹配用户侧最终 ASR,失败/未知挂断不自动重试。`ApprovedOriginator` 经生成的 Unary gRPC Stub 交付完整快照,SIP 全量 revision 不吻合即拒绝。 +- 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 或第二阶段外部验收。 diff --git a/internal/ai/current_opening_test.go b/internal/ai/current_opening_test.go new file mode 100644 index 0000000..bf81a91 --- /dev/null +++ b/internal/ai/current_opening_test.go @@ -0,0 +1,138 @@ +package ai + +import ( + "context" + "encoding/base64" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" +) + +func TestCurrentCallOpeningUsesApprovedTTSOnce(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + var calls atomic.Int32 + text := make(chan string, 1) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls.Add(1) + 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 + } + text <- payload.Params.Text + _, _ = fmt.Fprintf(w, `{"code":0,"data":%q}`+"\n", base64.StdEncoding.EncodeToString([]byte{1, 0, 2, 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) + } + call, err := NewCurrentCall(bound, func(context.Context) error { return nil }) + if err != nil { + t.Fatal(err) + } + audio, err := call.Open(context.Background()) + if err != nil || len(audio) != 4 || <-text != bound.Opening { + t.Fatalf("approved opening was not synthesized once: len=%d err=%v", len(audio), err) + } + if _, err := call.Open(context.Background()); err == nil || !strings.Contains(err.Error(), "already") || calls.Load() != 1 { + t.Fatalf("opening cannot be replayed: requests=%d err=%v", calls.Load(), err) + } +} + +func TestCurrentCallOpeningFailureIsVisibleAndNeverRetried(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + var calls atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls.Add(1) + http.Error(w, "mock provider unavailable", http.StatusServiceUnavailable) + })) + 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) + } + call, err := NewCurrentCall(bound, func(context.Context) error { return nil }) + if err != nil { + t.Fatal(err) + } + if _, err := call.Open(context.Background()); err == nil { + t.Fatal("provider failure must be visible") + } + if _, err := call.Open(context.Background()); err == nil || calls.Load() != 1 { + t.Fatalf("uncertain opening must not replay: requests=%d err=%v", calls.Load(), err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if _, err := call.RunTurn(ctx, []byte{1, 0}); err == nil || !strings.Contains(err.Error(), "opening") { + t.Fatalf("an unsuccessful opening cannot be silently skipped: %v", err) + } +} + +func TestCurrentCallASROnlyNeverRequestsOpeningTTS(t *testing.T) { + task, providers := currentFixture(t, "asr_only") + bound, err := BindCurrent(task, providers) + if err != nil { + t.Fatal(err) + } + call, err := NewCurrentCall(bound, func(context.Context) error { return nil }) + if err != nil { + t.Fatal(err) + } + if _, err := call.Open(context.Background()); err == nil { + t.Fatal("ASR-only call cannot synthesize an opening") + } +} + +func TestCurrentCallOptionalEmptyOpeningMakesNoTTSRequest(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + task = changeCurrentAgent(t, task, func(agent map[string]any) { + agent["conversation"].(map[string]any)["opening"] = "" + }) + bound, err := BindCurrent(task, providers) + if err != nil { + t.Fatal(err) + } + call, err := NewCurrentCall(bound, func(context.Context) error { return nil }) + if err != nil { + t.Fatal(err) + } + if audio, err := call.Open(context.Background()); err != nil || len(audio) != 0 { + t.Fatalf("an approved empty opening must not call TTS: len=%d err=%v", len(audio), err) + } + if _, err := call.Open(context.Background()); err == nil { + t.Fatal("opening stage cannot run twice even if no audio was configured") + } +} + +func TestCurrentCallRequiresApprovedOpeningBeforeFullAITurn(t *testing.T) { + task, providers := currentFixture(t, "full_ai") + bound, err := BindCurrent(task, providers) + if err != nil { + t.Fatal(err) + } + call, err := NewCurrentCall(bound, func(context.Context) error { return nil }) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if _, err := call.RunTurn(ctx, []byte{1, 0}); err == nil || !strings.Contains(err.Error(), "opening") { + t.Fatalf("a full-AI turn cannot skip its configured opening: %v", err) + } +} diff --git a/internal/ai/current_pipeline.go b/internal/ai/current_pipeline.go index f546285..9d7ec6e 100644 --- a/internal/ai/current_pipeline.go +++ b/internal/ai/current_pipeline.go @@ -6,6 +6,7 @@ import ( "fmt" "net/url" "strings" + "sync" "time" "github.com/GizClaw/doubao-speech-go" @@ -18,8 +19,11 @@ import ( // CurrentCall owns the keyword action for exactly one call. A failed or // uncertain hangup is never attempted again for that call. type CurrentCall struct { - bound CurrentBound - keyword *KeywordHangup + bound CurrentBound + keyword *KeywordHangup + mu sync.Mutex + openingRequested bool + openingReady bool } func NewCurrentCall(bound CurrentBound, hangup func(context.Context) error) (*CurrentCall, error) { @@ -30,6 +34,35 @@ func NewCurrentCall(bound CurrentBound, hangup func(context.Context) error) (*Cu return &CurrentCall{bound: bound, keyword: keyword}, nil } +// Open synthesizes the approved opening at most once per call. An ambiguous +// synthesis result is visible to the caller and never replayed automatically. +func (c *CurrentCall) Open(ctx context.Context) ([]byte, error) { + if c == nil || c.bound.Mode != "full_ai" { + return nil, errors.New("ASR-only call has no opening TTS") + } + c.mu.Lock() + if c.openingRequested { + c.mu.Unlock() + return nil, errors.New("opening already requested; outcome unknown") + } + c.openingRequested = true + c.mu.Unlock() + if c.bound.Opening == "" { + c.mu.Lock() + c.openingReady = true + c.mu.Unlock() + return nil, nil + } + audio, err := c.bound.Synthesize(ctx, c.bound.Opening) + if err != nil { + return nil, err + } + c.mu.Lock() + c.openingReady = true + c.mu.Unlock() + return audio, nil +} + func (c *CurrentCall) HandleFinalASR(ctx context.Context, segment ASRSegment) (bool, error) { if c == nil { return false, errors.New("AI call is unavailable") @@ -43,6 +76,14 @@ func (c *CurrentCall) RunTurn(ctx context.Context, pcm16 []byte) (TurnResult, er if c == nil { return TurnResult{}, errors.New("AI call is unavailable") } + if c.bound.Mode == "full_ai" && c.bound.Opening != "" { + c.mu.Lock() + ready := c.openingReady + c.mu.Unlock() + if !ready { + return TurnResult{}, errors.New("approved opening was not completed") + } + } text, err := c.bound.Recognize(ctx, pcm16) if err != nil { return TurnResult{}, fmt.Errorf("ASR failed: %w", err) diff --git a/internal/rpc/approved_full_ai_integration_test.go b/internal/rpc/approved_full_ai_integration_test.go index bde63f3..c380566 100644 --- a/internal/rpc/approved_full_ai_integration_test.go +++ b/internal/rpc/approved_full_ai_integration_test.go @@ -39,7 +39,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T req.TaskId = task.TaskID req.SourceEventId, req.CallId = "event-full", "event-full" llmObserved := make(chan approvedSDKObservation, 1) - ttsObserved := make(chan approvedSDKObservation, 1) + ttsObserved := make(chan approvedSDKObservation, 2) mockSDKs := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { var body map[string]any if err := json.NewDecoder(r.Body).Decode(&body); err != nil { @@ -89,6 +89,13 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T if err != nil { return err } + openingAudio, err := call.Open(ctx) + if err != nil { + return fmt.Errorf("synthesize approved opening: %w", err) + } + if len(openingAudio) != 4 { + t.Fatal("Agent did not synthesize the approved opening audio") + } for _, segment := range []ai.ASRSegment{ {Source: "user", Text: "不用了", Final: false}, {Source: "assistant", Text: "不用了", Final: true}, @@ -136,7 +143,8 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T if err := orig.Originate(context.Background(), spec); err != nil || calls != 1 || hangups != 1 { t.Fatalf("full-AI isolated Agent call failed: err=%v calls=%d hangups=%d", err, calls, hangups) } - llmCall, ttsCall := <-llmObserved, <-ttsObserved + llmCall := <-llmObserved + openingCall, ttsCall := <-ttsObserved, <-ttsObserved if llmCall.Authorization != "Bearer "+llm.Credential || llmCall.Body["model"] != "example-chat" || llmCall.Body["temperature"] != float64(0) || llmCall.Body["max_tokens"] != float64(256) { t.Fatal("Agent did not send approved LLM values through SDK") } @@ -144,7 +152,7 @@ func TestApprovedFullAIUsesBoundProviderSDKsAndOnlyFinalUserKeyword(t *testing.T t.Fatal("Agent did not use approved TTS resource and credential") } params := ttsCall.Body["req_params"].(map[string]any) - if params["speaker"] != "example-neutral" || params["text"] != "回答" || params["audio_params"].(map[string]any)["format"] != "pcm_s16le" { - t.Fatal("Agent did not send approved TTS voice/text/format") + if openingCall.Body["req_params"].(map[string]any)["text"] != "Example greeting" || params["speaker"] != "example-neutral" || params["text"] != "回答" || params["audio_params"].(map[string]any)["format"] != "pcm_s16le" { + t.Fatal("Agent did not send the approved opening then reply via TTS") } }