From fe33edb11c112dd270cbf0b7bbf6708e46b5e162 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 02:26:07 +0800 Subject: [PATCH] Build final call results from observed media and approved facts --- .../saas-dispatcher-implementation.md | 1 + internal/callflow/capture.go | 36 ++++-- internal/callflow/capture_test.go | 4 +- internal/callflow/flow.go | 25 ++-- internal/callflow/flow_test.go | 19 +++ internal/callflow/result_payload.go | 100 +++++++++++++++ internal/callflow/result_payload_test.go | 116 ++++++++++++++++++ 7 files changed, 278 insertions(+), 23 deletions(-) create mode 100644 internal/callflow/result_payload.go create mode 100644 internal/callflow/result_payload_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 65183c7..3337a09 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -67,6 +67,7 @@ - Dispatcher 录音事实 Unary RPC 隔离服务:`RequestRecordingUpload`、`ReportCallEnded`、`ReportCallResult` 均要求已配置本 D、核验 mTLS 指纹及当前 Agent 会话、数字租户和已保留的执行;复用官方 SDK 仅对原始录音签发固定 15 分钟授权,显式重申请仍用相同 bucket/object_key。结束事实可先于外呼响应而持久化原回执;录音结果核对 D 已存目标和 Agent 报告的成功 PUT,再与唯一结果 outbox 同事务提交。Mock 覆盖会话/租户拒绝、同资产重申请、上传前结果拒绝、坏 JSON、已上传与无录音结果及重复回报。隔离测试还通过本地双向 TLS 的 gRPC 实际传输:Agent `RecordingClient` 每次读取并克隆当前会话元数据,经受控 D 客户端领取授权、上报结束和唯一结果;未配置客户端明确拒绝。其分层测试不包含 OSS PUT;下述本地隔离链路另验证直传,主入口批准执行的媒体录音仍未接线。 - 内存录音隔离组件:`RecordingSession` 仅复制共享通话流程实际读到和成功发送的 16-kHz PCM16,`EncodeMonoWAV` 直接在内存生成有界单声道 WAV;空音频、奇数字节、超过上限及未成功发送的音频都不能伪造成可上传录音。单元与 race 测试未产生业务文件。批准执行入口尚未接入该组件,且 Mock 中观测到的帧不等于真实 Asterisk 通话的全量媒体验收。 - Agent 录音交付隔离组件:`RecordingDelivery` 先确认结束,再依照录音是否实际生成分别上报唯一空录音结果或请求原授权并直传内存 WAV;录音生成失败保留通话真实结果、空录音对象及明确原因,不虚构上传事实。隔离测试通过本地 HTTP PUT 和假 Dispatcher RPC 覆盖成功无业务文件、OSS 明确失败后私有文件保存、恢复写入失败、未知 PUT 隔离、重启重领原目标、上传已确认后只重发原结果。再次调用不会隐式重新 PUT;正常已确认上传但尚未被 D 持久收讫的跨进程间隙仍受 K16 边界约束。此处未连接主入口批准执行媒体或 MQ;下述隔离链路另测本地真正的 D gRPC/SQLite。 +- 最终结果隔离组件:`FinalResultPayload` 只从批准的任务身份、确认时间及实际收到的用户媒体窗口生成当前唯一结果;最终 ASR 文本与字面关键词拒联写入完整转写,未播放的助手候选内容不冒充转写,未知 SIP 响应码保留 `null`,无录音保留 `{}`。单元测试直接核验当前 MQ Schema、字段缺失/乱序时窗与无应答负例;开场白发送失败不得计入已播放音频。此组件尚未接入主入口,也未证明真实 RTP 转写时间精度。 - 本地隔离链路:`RecordingSession` 从实际读出的 Mock 媒体生成内存 WAV,经 Agent→Dispatcher 双向 TLS gRPC 确认结束、领取原始资产的签名授权,再对本地 HTTPS OSS Mock 单次 PUT;Dispatcher 将原上传事实与唯一最终结果 outbox 同事务提交。测试确认一次 PUT、一次结果、无业务文件,使用的是官方 SDK 生成的路径,OSS Mock 不验证真实服务商签名;这不是主入口执行、RabbitMQ 投递或外部验收。 - 已验证:`go test ./... -count=1`、`go test -race ./... -count=1`、`go vet ./...`、`go build ./...`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`git diff --check`。尚未完成批准执行媒体录音到 Agent 交付组件的主入口接线、主入口 Agent↔Dispatcher 实际会话与录音执行接线(隔离 mTLS gRPC 已测)、真实执行时上传事实/最终结果交付及 MQ/端到端验收,不能宣称 P06 通过。 diff --git a/internal/callflow/capture.go b/internal/callflow/capture.go index 35f158d..fdc0626 100644 --- a/internal/callflow/capture.go +++ b/internal/callflow/capture.go @@ -20,7 +20,12 @@ type CaptureConfig struct { MaxPendingAudioChunks int } -func captureTurn(ctx context.Context, session MediaSession, cfg CaptureConfig) ([]byte, error) { +type capturedTurn struct { + PCM []byte + Window CaptureWindow +} + +func captureTurn(ctx context.Context, session MediaSession, cfg CaptureConfig) (capturedTurn, error) { if cfg.FirstSpeechTimeout <= 0 { cfg.FirstSpeechTimeout = 5 * time.Second } @@ -35,15 +40,16 @@ func captureTurn(ctx context.Context, session MediaSession, cfg CaptureConfig) ( maxDeadline := startedAt.Add(cfg.MaxDuration) started := cfg.VoiceThreshold <= 0 lastVoice := time.Time{} + var firstFrameAt, lastFrameAt time.Time var frames []byte for { if err := ctx.Err(); err != nil { - return nil, err + return capturedTurn{}, err } now := time.Now() if !started && !now.Before(firstDeadline) { - return nil, errors.New("no speech detected before capture timeout") + return capturedTurn{}, errors.New("no speech detected before capture timeout") } if !maxDeadline.After(now) { break @@ -62,7 +68,7 @@ func captureTurn(ctx context.Context, session MediaSession, cfg CaptureConfig) ( payload, err := session.ReadPayload(readCtx) cancel() if err := ctx.Err(); err != nil { - return nil, err + return capturedTurn{}, err } if err != nil { if errors.Is(err, context.DeadlineExceeded) { @@ -74,22 +80,26 @@ func captureTurn(ctx context.Context, session MediaSession, cfg CaptureConfig) ( if errors.Is(err, context.Canceled) && ctx.Err() == nil { continue } - return nil, err + return capturedTurn{}, err } if len(payload) == 0 { continue } + receivedAt := time.Now() if cfg.VoiceThreshold > 0 { if pcm16VoiceLevel(payload) >= cfg.VoiceThreshold { started = true - lastVoice = time.Now() - } - if started { - frames = append(frames, payload...) + lastVoice = receivedAt } } else { started = true - lastVoice = time.Now() + lastVoice = receivedAt + } + if started { + if firstFrameAt.IsZero() { + firstFrameAt = receivedAt + } + lastFrameAt = receivedAt frames = append(frames, payload...) } if started && cfg.EndSilence > 0 && !lastVoice.IsZero() && !time.Now().Before(lastVoice.Add(cfg.EndSilence)) { @@ -97,12 +107,12 @@ func captureTurn(ctx context.Context, session MediaSession, cfg CaptureConfig) ( } } if err := ctx.Err(); err != nil { - return nil, err + return capturedTurn{}, err } if len(frames) == 0 { - return nil, errors.New("captured audio is empty") + return capturedTurn{}, errors.New("captured audio is empty") } - return frames, nil + return capturedTurn{PCM: frames, Window: CaptureWindow{StartedAt: firstFrameAt, EndedAt: lastFrameAt}}, nil } func pcm16VoiceLevel(pcm []byte) int { diff --git a/internal/callflow/capture_test.go b/internal/callflow/capture_test.go index 487b682..630cde9 100644 --- a/internal/callflow/capture_test.go +++ b/internal/callflow/capture_test.go @@ -23,7 +23,7 @@ func TestCaptureTurnWaitsForSpeechAndStopsOnSilence(t *testing.T) { if err != nil { t.Fatal(err) } - if len(captured) < len(voice) || len(captured) >= len(silence)+len(voice)+len(silence) { - t.Fatalf("captured=%d want speech and bounded trailing silence", len(captured)) + if len(captured.PCM) < len(voice) || len(captured.PCM) >= len(silence)+len(voice)+len(silence) || captured.Window.StartedAt.IsZero() || captured.Window.EndedAt.Before(captured.Window.StartedAt) { + t.Fatalf("capture lost voice, duration boundary or observed times: pcm=%d window=%+v", len(captured.PCM), captured.Window) } } diff --git a/internal/callflow/flow.go b/internal/callflow/flow.go index f2d7658..c00a14c 100644 --- a/internal/callflow/flow.go +++ b/internal/callflow/flow.go @@ -19,13 +19,20 @@ type MediaSession interface { Stats() media.RTPStats } +// CaptureWindow bounds the actual receipt of one user-side audio turn. +type CaptureWindow struct { + StartedAt time.Time + EndedAt time.Time +} + type Result struct { - Turn ai.TurnResult - Turns []ai.TurnResult - Inbound []byte - InboundTurns [][]byte - OutboundTurns [][]byte - RTP media.RTPStats + Turn ai.TurnResult + Turns []ai.TurnResult + Inbound []byte + InboundTurns [][]byte + CaptureWindows []CaptureWindow + OutboundTurns [][]byte + RTP media.RTPStats } func Execute(ctx context.Context, session MediaSession, pipeline ai.Pipeline, snapshot ai.Snapshot, opening string, turnWindow time.Duration) (Result, error) { @@ -69,19 +76,21 @@ func executeFlow(ctx context.Context, session MediaSession, mode ai.Mode, open f 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 } + result.OutboundTurns = append(result.OutboundTurns, clonePCM(openingPCM)) } } for turnIndex := 0; turnIndex < maxTurns; turnIndex++ { - inbound, err := captureTurn(ctx, session, capture) + captured, err := captureTurn(ctx, session, capture) if err != nil { return result, fmt.Errorf("capturing RTP turn %d after opening/reply prompt: %w", turnIndex+1, err) } + inbound := captured.PCM result.Inbound = inbound result.InboundTurns = append(result.InboundTurns, clonePCM(inbound)) + result.CaptureWindows = append(result.CaptureWindows, captured.Window) if len(inbound) < 3200 { return result, fmt.Errorf("captured audio turn %d is too short", turnIndex+1) } diff --git a/internal/callflow/flow_test.go b/internal/callflow/flow_test.go index 9c7092d..df75101 100644 --- a/internal/callflow/flow_test.go +++ b/internal/callflow/flow_test.go @@ -93,6 +93,20 @@ func (invalidCallPipeline) RunTurn(context.Context, ai.Snapshot, []byte) (ai.Tur 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 { @@ -130,14 +144,19 @@ func TestASROnlySkipsOpeningLLMAndTTS(t *testing.T) { } 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) } diff --git a/internal/callflow/result_payload.go b/internal/callflow/result_payload.go new file mode 100644 index 0000000..5d49f29 --- /dev/null +++ b/internal/callflow/result_payload.go @@ -0,0 +1,100 @@ +package callflow + +import ( + "encoding/json" + "errors" + "fmt" + "strings" + "time" +) + +// FinalCallFacts are observed or approved call facts, never defaults inferred +// from the Agent SDK. The reason code is nil if no actual SIP code is known. +type FinalCallFacts struct { + TaskID string + CallerProfileID string + Callee string + TrunkID string + StartedAt time.Time + EndedAt time.Time + Outcome string + ReasonCode *int + ReasonMessage string +} + +// FinalResultPayload creates the one current MQ result body. Only final +// user-side ASR text with observed capture timing is included; unplayed LLM +// replies, interim ASR and hypothetical recordings are never transcribed. +func FinalResultPayload(facts FinalCallFacts, call Result) ([]byte, error) { + if strings.TrimSpace(facts.TaskID) == "" || strings.TrimSpace(facts.CallerProfileID) == "" || strings.TrimSpace(facts.Callee) == "" || strings.TrimSpace(facts.TrunkID) == "" || strings.TrimSpace(facts.ReasonMessage) == "" { + return nil, errors.New("final call result is missing approved identity or observed reason") + } + if facts.StartedAt.IsZero() || facts.EndedAt.IsZero() || facts.EndedAt.Before(facts.StartedAt) { + return nil, errors.New("final call result has no valid observed time range") + } + if facts.Outcome != "answered" && facts.Outcome != "no_answer" && facts.Outcome != "failed" { + return nil, errors.New("final call result has no approved outcome") + } + if facts.Outcome == "no_answer" && len(call.Turns) != 0 { + return nil, errors.New("no-answer call cannot contain completed user turns") + } + if facts.ReasonCode != nil && (*facts.ReasonCode < 100 || *facts.ReasonCode > 699) { + return nil, errors.New("final call result has an invalid SIP response code") + } + if len(call.CaptureWindows) < len(call.Turns) { + return nil, errors.New("final ASR text has no corresponding observed media window") + } + type transcriptSegment struct { + SegmentID string `json:"segment_id"` + TurnID string `json:"turn_id"` + Role string `json:"role"` + StartMS int64 `json:"start_ms"` + EndMS int64 `json:"end_ms"` + Text string `json:"text"` + } + segments := make([]transcriptSegment, 0, len(call.Turns)) + optOut := false + callDuration := facts.EndedAt.Sub(facts.StartedAt).Milliseconds() + for i, turn := range call.Turns { + window := call.CaptureWindows[i] + if window.StartedAt.Before(facts.StartedAt) || window.EndedAt.Before(window.StartedAt) || window.EndedAt.After(facts.EndedAt) { + return nil, fmt.Errorf("turn %d has no valid observed capture times", i+1) + } + start, end := window.StartedAt.Sub(facts.StartedAt).Milliseconds(), window.EndedAt.Sub(facts.StartedAt).Milliseconds() + if start < 0 || end < start || end > callDuration { + return nil, fmt.Errorf("turn %d capture is outside confirmed call duration", i+1) + } + if turn.EndedByKeyword { + if strings.TrimSpace(turn.Transcript) == "" { + return nil, fmt.Errorf("turn %d claims opt-out without final user ASR text", i+1) + } + optOut = true + } + if turn.Transcript != "" { + segments = append(segments, transcriptSegment{ + SegmentID: fmt.Sprintf("segment-%d", len(segments)+1), TurnID: fmt.Sprintf("turn-%d", i+1), + Role: "user", StartMS: start, EndMS: end, Text: turn.Transcript, + }) + } + } + return json.Marshal(struct { + TaskID string `json:"task_id"` + CallerProfileID string `json:"caller_profile_id"` + Callee string `json:"callee"` + TrunkID string `json:"trunk_id"` + StartedAt string `json:"started_at"` + EndedAt string `json:"ended_at"` + DurationMS int64 `json:"duration_ms"` + Outcome string `json:"outcome"` + ReasonCode *int `json:"reason_code"` + ReasonMessage string `json:"reason_message"` + Transcript []transcriptSegment `json:"transcript"` + OptOut bool `json:"opt_out"` + Recording map[string]any `json:"recording"` + }{ + TaskID: facts.TaskID, CallerProfileID: facts.CallerProfileID, Callee: facts.Callee, TrunkID: facts.TrunkID, + StartedAt: facts.StartedAt.UTC().Format(time.RFC3339Nano), EndedAt: facts.EndedAt.UTC().Format(time.RFC3339Nano), + DurationMS: callDuration, Outcome: facts.Outcome, ReasonCode: facts.ReasonCode, + ReasonMessage: facts.ReasonMessage, Transcript: segments, OptOut: optOut, Recording: map[string]any{}, + }) +} diff --git a/internal/callflow/result_payload_test.go b/internal/callflow/result_payload_test.go new file mode 100644 index 0000000..b737158 --- /dev/null +++ b/internal/callflow/result_payload_test.go @@ -0,0 +1,116 @@ +package callflow + +import ( + "encoding/json" + "testing" + "time" + + "git.ipao.vip/rogee/go-sip/internal/ai" + "git.ipao.vip/rogee/go-sip/internal/contract" +) + +func approvedResultFixture() (FinalCallFacts, Result) { + start := time.Date(2026, 9, 21, 1, 0, 0, 0, time.UTC) + facts := FinalCallFacts{ + TaskID: "task-asr", CallerProfileID: "caller-profile-mock", Callee: "15003164745", TrunkID: "trunk-mock", + StartedAt: start, EndedAt: start.Add(3 * time.Second), Outcome: "answered", ReasonMessage: "completed in isolated mock", + } + result := Result{ + Turns: []ai.TurnResult{ + {Transcript: "用户真实识别文本", Reply: "未实际播放的助手候选内容"}, + {Transcript: "请不要联系", EndedByKeyword: true}, + }, + CaptureWindows: []CaptureWindow{ + {StartedAt: start.Add(100 * time.Millisecond), EndedAt: start.Add(1200 * time.Millisecond)}, + {StartedAt: start.Add(1600 * time.Millisecond), EndedAt: start.Add(2250 * time.Millisecond)}, + }, + } + return facts, result +} + +func TestFinalResultPayloadUsesOnlyObservedUserASRAndMediaTimes(t *testing.T) { + facts, result := approvedResultFixture() + payload, err := FinalResultPayload(facts, result) + if err != nil { + t.Fatal(err) + } + message, err := json.Marshal(map[string]any{ + "event_id": "call-result-test", "event_type": "call.execute.result", + "dispatcher_id": "c046b893-8628-4589-ae50-619d049248a6", "tenant_id": 1001, + "issued_at": "2026-09-21T01:00:04Z", "payload": json.RawMessage(payload), + }) + if err != nil { + t.Fatal(err) + } + if err := contract.ValidateCurrent("mq", message); err != nil { + t.Fatalf("locally produced result violates the sole current MQ schema: %v", err) + } + var parsed struct { + TaskID string `json:"task_id"` + CallerProfileID string `json:"caller_profile_id"` + DurationMS int64 `json:"duration_ms"` + Outcome string `json:"outcome"` + ReasonCode *int `json:"reason_code"` + OptOut bool `json:"opt_out"` + Recording map[string]any `json:"recording"` + Transcript []struct { + SegmentID string `json:"segment_id"` + TurnID string `json:"turn_id"` + Role string `json:"role"` + StartMS int64 `json:"start_ms"` + EndMS int64 `json:"end_ms"` + Text string `json:"text"` + } `json:"transcript"` + } + if err := json.Unmarshal(payload, &parsed); err != nil { + t.Fatal(err) + } + if parsed.TaskID != facts.TaskID || parsed.CallerProfileID != facts.CallerProfileID || parsed.DurationMS != 3000 || parsed.Outcome != "answered" || parsed.ReasonCode != nil || !parsed.OptOut || len(parsed.Recording) != 0 || len(parsed.Transcript) != 2 { + t.Fatalf("result changed approved facts or invented a recording: %+v", parsed) + } + if first, second := parsed.Transcript[0], parsed.Transcript[1]; first.SegmentID != "segment-1" || first.TurnID != "turn-1" || first.Role != "user" || first.StartMS != 100 || first.EndMS != 1200 || first.Text != result.Turns[0].Transcript || second.SegmentID != "segment-2" || second.TurnID != "turn-2" || second.StartMS != 1600 || second.EndMS != 2250 || second.Text != result.Turns[1].Transcript { + t.Fatalf("ASR text or observed media timing was lost: %+v", parsed.Transcript) + } +} + +func TestFinalResultPayloadRejectsMissingOrInventedCallFacts(t *testing.T) { + for _, tc := range []struct { + name string + change func(*FinalCallFacts, *Result) + }{ + {"missing_approved_caller", func(f *FinalCallFacts, _ *Result) { f.CallerProfileID = "" }}, + {"unapproved_outcome", func(f *FinalCallFacts, _ *Result) { f.Outcome = "new_state" }}, + {"no_sip_reason", func(f *FinalCallFacts, _ *Result) { f.ReasonMessage = "" }}, + {"ended_before_started", func(f *FinalCallFacts, _ *Result) { f.EndedAt = f.StartedAt.Add(-time.Second) }}, + {"missing_media_window", func(_ *FinalCallFacts, r *Result) { r.CaptureWindows = nil }}, + {"reversed_media_window", func(_ *FinalCallFacts, r *Result) { + r.CaptureWindows[0].EndedAt = r.CaptureWindows[0].StartedAt.Add(-time.Millisecond) + }}, + {"window_outside_confirmed_call", func(f *FinalCallFacts, r *Result) { r.CaptureWindows[0].EndedAt = f.EndedAt.Add(time.Millisecond) }}, + } { + t.Run(tc.name, func(t *testing.T) { + facts, result := approvedResultFixture() + tc.change(&facts, &result) + if payload, err := FinalResultPayload(facts, result); err == nil || len(payload) != 0 { + t.Fatalf("unconfirmed fact became a final result: payload_size=%d err=%v", len(payload), err) + } + }) + } +} + +func TestFinalResultPayloadNoAnswerHasNoFabricatedTranscriptOrAsset(t *testing.T) { + facts, _ := approvedResultFixture() + facts.Outcome, facts.ReasonMessage = "no_answer", "SIP no answer" + payload, err := FinalResultPayload(facts, Result{}) + if err != nil { + t.Fatal(err) + } + var result struct { + Transcript []any `json:"transcript"` + OptOut bool `json:"opt_out"` + Recording map[string]any `json:"recording"` + } + if err := json.Unmarshal(payload, &result); err != nil || len(result.Transcript) != 0 || result.OptOut || len(result.Recording) != 0 { + t.Fatalf("no-answer outcome invented media facts: %+v err=%v", result, err) + } +}