diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index a9eecbe..5abc223 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -76,6 +76,7 @@ - Agent 录音交付隔离组件:`RecordingDelivery` 先确认结束,再依照录音是否实际生成分别上报唯一空录音结果或请求原授权并直传内存 WAV;录音生成失败保留通话真实结果、空录音对象及明确原因,不虚构上传事实。隔离测试通过本地 HTTP PUT 和假 Dispatcher RPC 覆盖成功无业务文件、OSS 明确失败后私有文件保存、恢复写入失败、未知 PUT 隔离、重启重领原目标、上传已确认后只重发原结果。再次调用不会隐式重新 PUT;正常已确认上传但尚未被 D 持久收讫的跨进程间隙仍受 K16 边界约束。此处未连接主入口批准执行媒体或 MQ;下述隔离链路另测本地真正的 D gRPC/SQLite。 - 最终结果隔离组件:`ApprovedExecution` 已保留获批任务的 `caller_profile_id`;`FinalResultPayload` 只从批准的任务身份、确认时间及实际收到的用户媒体窗口生成当前唯一结果;最终 ASR 文本与字面关键词拒联写入完整转写,未播放的助手候选内容不冒充转写,未知 SIP 响应码保留 `null`,无录音保留 `{}`。单元测试直接核验当前 MQ Schema、字段缺失/乱序时窗与无应答负例;开场白发送失败不得计入已播放音频。此组件尚未接入主入口,也未证明真实 RTP 转写时间精度。 - 本地隔离链路:`RunApprovedCall` 显式使用合成 `ApprovedMockPipeline`,让共享批准通话流程按 ASR-only 限制消费实际读出的 PCM Mock 媒体帧;`RecordingSession` 生成内存 WAV,`FinalResultPayload` 仅用最终模拟识别及实测采集时窗构造结果。Agent 经双向 TLS gRPC 向 Dispatcher 确认结束、领取原始资产签名授权,再对本地 HTTPS OSS Mock 单次 PUT;Dispatcher 将原上传事实与唯一最终结果 outbox 同事务提交。测试确认一次 PUT、一次结果、无业务文件,使用官方 SDK 生成的路径;Mock ASR 与 OSS 服务不验证真实供应商协议或签名。这不是主入口执行、RabbitMQ 投递或外部验收。 +- `ApprovedRecordedMockCall` 将显式批准的 ASR-only AI 快照、合成 PCM 媒体、每通话独立脚本、实际采集的有界 WAV、用户最终识别时窗和 `RecordingDelivery` 组合到一个可供 Agent 工人调用的 Mock runner;隔离测试验证先结束事实、一次本地 OSS PUT、唯一最终结果及正常路径无业务文件。无实际读入媒体时明确报告失败并以空录音、生成失败原因收口,不伪造上传或转写;缺失每通话交付器在启动媒体前报错。此 runner 尚未接到主 CLI 的 `ApprovedCallWorker`,测试 OSS/RPC 均为隔离替身,不能代表真实线路、供应商或 MQ 投递。 - 已验证:`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/rpc/approved_recorded_mock.go b/internal/rpc/approved_recorded_mock.go new file mode 100644 index 0000000..0a48b80 --- /dev/null +++ b/internal/rpc/approved_recorded_mock.go @@ -0,0 +1,86 @@ +package rpc + +import ( + "context" + "errors" + "strings" + "time" + + "git.ipao.vip/rogee/go-sip/internal/agent" + "git.ipao.vip/rogee/go-sip/internal/ai" + "git.ipao.vip/rogee/go-sip/internal/callflow" +) + +var ErrApprovedMockDeliveryUnavailable = errors.New("approved Mock call requires a per-call recording delivery") + +// ApprovedRecordedMockCall joins the approved callflow with bounded in-memory +// capture and one final delivery. The synthetic input and script are explicit +// test fixtures, not SaaS AI settings or evidence of a real SIP/provider call. +// A new instance and RecordingDelivery belong to each call. +type ApprovedRecordedMockCall struct { + InboundPCM16 []byte + Script ai.ApprovedMockScript + MaxWAVBytes int64 + ExpectedRecording bool + Outcome string + ReasonMessage string + ReportTimeout time.Duration + Delivery *agent.RecordingDelivery +} + +// Run ends the media flow before releasing its caller's task registration. +// Reporting uses a separately bounded context so a task hangup cannot silently +// suppress its confirmed end or the single final result. +func (r *ApprovedRecordedMockCall) Run(ctx context.Context, approved ApprovedExecution) error { + if r == nil || ctx == nil || r.Delivery == nil || r.ReportTimeout <= 0 { + return ErrApprovedMockDeliveryUnavailable + } + if approved.DispatcherID == "" || approved.TenantID <= 0 || approved.SourceEventID == "" || approved.CallID != approved.SourceEventID || + approved.TaskID == "" || approved.CallerProfileID == "" || approved.Callee == "" || approved.SelectedTrunkID == "" || r.MaxWAVBytes <= 44 || strings.TrimSpace(r.ReasonMessage) == "" { + return errors.New("approved Mock call identity, media limit or observed reason is incomplete") + } + if r.Outcome != "answered" && r.Outcome != "no_answer" && r.Outcome != "failed" { + return errors.New("approved Mock call outcome is invalid") + } + if r.ExpectedRecording && (r.Delivery.Recovery == nil || strings.TrimSpace(r.Delivery.Recovery.Root) == "") { + return ErrApprovedMockDeliveryUnavailable + } + media, err := callflow.NewRecordingSession(callflow.NewMemorySession(r.InboundPCM16), r.MaxWAVBytes) + if err != nil { + return err + } + // The isolated memory session has no external SIP channel to hang up. + hangup := func(context.Context) error { return nil } + pipeline, err := ai.NewApprovedMockPipeline(approved.AI, r.Script, hangup) + if err != nil { + return err + } + startedAt := time.Now().UTC() + observed, runErr := RunApprovedCall(ctx, approved, media, hangup, pipeline) + endedAt := time.Now().UTC() + outcome, reason := r.Outcome, r.ReasonMessage + if runErr != nil { + outcome, reason = "failed", "isolated Mock media did not complete" + } + payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{ + TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID, + Callee: approved.Callee, TrunkID: approved.SelectedTrunkID, + StartedAt: startedAt, EndedAt: endedAt, Outcome: outcome, ReasonMessage: reason, + }, observed) + if err != nil { + return errors.Join(runErr, err) + } + completed := agent.CompletedRecording{ResultPayload: payload, Expected: r.ExpectedRecording} + if r.ExpectedRecording { + wav, durationMS, captureErr := media.WAV() + completed.CaptureError = captureErr + if captureErr == nil { + completed.RecordingID = "recording-" + approved.SourceEventID + completed.UploadID = "upload-" + approved.SourceEventID + completed.WAV, completed.DurationMS = wav, durationMS + } + } + reportCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), r.ReportTimeout) + defer cancel() + return errors.Join(runErr, r.Delivery.Complete(reportCtx, completed)) +} diff --git a/internal/rpc/approved_recorded_mock_test.go b/internal/rpc/approved_recorded_mock_test.go new file mode 100644 index 0000000..0fdf8cb --- /dev/null +++ b/internal/rpc/approved_recorded_mock_test.go @@ -0,0 +1,167 @@ +package rpc + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "io" + "net/http" + "net/http/httptest" + "os" + "strings" + "sync/atomic" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/agent" + "git.ipao.vip/rogee/go-sip/internal/ai" + "git.ipao.vip/rogee/go-sip/internal/media" + "google.golang.org/grpc" +) + +type recordedMockRPC struct { + agentpb.AgentControlServiceClient + targetURL string + calls []string + asset *agentpb.AssetDescriptor + result []byte + proof *agentpb.UploadObservation +} + +func (f *recordedMockRPC) RequestRecordingUpload(_ context.Context, req *agentpb.RequestRecordingUploadRequest, _ ...grpc.CallOption) (*agentpb.RequestRecordingUploadResponse, error) { + f.calls = append(f.calls, "grant") + f.asset = req.Asset + return &agentpb.RequestRecordingUploadResponse{Grant: &agentpb.UploadGrant{ + UploadId: req.UploadId, Bucket: "mock-bucket", ObjectKey: "approved/fixture.wav", TargetUrl: f.targetURL, + ExpiresAtUnixMs: time.Now().Add(15 * time.Minute).UnixMilli(), MaxBytes: req.Asset.SizeBytes, + RequiredChecksumSha256: req.Asset.ChecksumSha256, + }}, nil +} + +func (f *recordedMockRPC) ReportCallEnded(_ context.Context, req *agentpb.ReportCallEndedRequest, _ ...grpc.CallOption) (*agentpb.ReportCallEndedResponse, error) { + f.calls = append(f.calls, "end") + return &agentpb.ReportCallEndedResponse{Receipt: &agentpb.OperationReceipt{ + FactId: req.SourceEventId, Result: agentpb.ResultCode_RESULT_CODE_APPLIED, + }}, nil +} + +func (f *recordedMockRPC) ReportCallResult(_ context.Context, req *agentpb.ReportCallResultRequest, _ ...grpc.CallOption) (*agentpb.ReportCallResultResponse, error) { + f.calls = append(f.calls, "result") + f.result = bytes.Clone(req.ResultPayloadJson) + f.proof = req.Upload + return &agentpb.ReportCallResultResponse{Receipt: &agentpb.OperationReceipt{ + FactId: "result-fixture", Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED, + }}, nil +} + +func recordedMockFixture(t *testing.T, pcm []byte, maxDuration time.Duration) (*ApprovedRecordedMockCall, *recordedMockRPC, *atomic.Int32, string, ApprovedExecution) { + t.Helper() + var puts atomic.Int32 + var sent []byte + localOSS := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + puts.Add(1) + body, err := io.ReadAll(io.LimitReader(r.Body, 4097)) + if err != nil || r.Method != http.MethodPut || r.URL.Path != "/object" { + t.Errorf("incorrect isolated OSS PUT: method=%s path=%s err=%v", r.Method, r.URL.Path, err) + w.WriteHeader(http.StatusBadRequest) + return + } + sent = bytes.Clone(body) + w.WriteHeader(http.StatusCreated) + })) + t.Cleanup(localOSS.Close) + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + stub := &recordedMockRPC{targetURL: localOSS.URL + "/object"} + approved := ApprovedExecution{ + DispatcherID: "c046b893-8628-4589-ae50-619d049248a6", TenantID: 42, + TaskID: "task-asr", CallerProfileID: "caller-mock", SourceEventID: "event-fixture", CallID: "event-fixture", + Callee: "15003164745", SelectedTrunkID: "trunk-mock", MaxCallDuration: maxDuration, + DialBefore: time.Now().Add(time.Minute), AI: ai.CurrentBound{ + Mode: "asr_only", Conversation: ai.CurrentConversation{SilenceTimeout: 100 * time.Millisecond, MaxDuration: maxDuration}, + }, + } + delivery := &agent.RecordingDelivery{ + Call: agent.RecordingClient{ + Client: stub, DispatcherID: approved.DispatcherID, TenantID: approved.TenantID, SourceEventID: approved.SourceEventID, + Session: func(context.Context) (*agentpb.RequestMeta, error) { + return &agentpb.RequestMeta{AgentId: "agent-mock", CellId: "cell-mock", BootId: "boot-mock", DispatcherEpoch: "epoch-mock", SessionGeneration: 1}, nil + }, + }, + Recovery: &agent.RecordingRecovery{Root: root, Upload: agent.UploadClient{AllowInsecureHTTP: true}}, + } + runner := &ApprovedRecordedMockCall{ + InboundPCM16: pcm, Script: ai.ApprovedMockScript{Turns: []ai.ApprovedMockTurn{{Transcript: "synthetic final ASR fixture"}}}, + MaxWAVBytes: 4096, ExpectedRecording: true, Outcome: "answered", ReasonMessage: "isolated Mock answered", + ReportTimeout: 5 * time.Second, Delivery: delivery, + } + // The test inspects the actual bytes accepted by the isolated OSS endpoint. + t.Cleanup(func() { + if puts.Load() > 0 && len(sent) == 0 { + t.Error("isolated OSS accepted an empty recording") + } + }) + return runner, stub, &puts, root, approved +} + +func TestApprovedRecordedMockCallDeliversObservedASRAndInMemoryWAV(t *testing.T) { + pcm := bytes.Repeat([]byte{1, 0}, 1600) + runner, stub, puts, root, approved := recordedMockFixture(t, pcm, time.Second) + if err := runner.Run(context.Background(), approved); err != nil { + t.Fatal(err) + } + if strings.Join(stub.calls, ",") != "end,grant,result" || puts.Load() != 1 || stub.asset == nil || stub.proof == nil { + t.Fatalf("observed media did not yield one confirmed result: calls=%v puts=%d", stub.calls, puts.Load()) + } + wav, duration, err := media.EncodeMonoWAV(pcm, 4096) + if err != nil || stub.asset.SizeBytes != int64(len(wav)) || stub.asset.DurationMs != duration || stub.asset.AssetId != "recording-"+approved.SourceEventID || stub.proof.RecordingId != stub.asset.AssetId || stub.proof.UploadId != "upload-"+approved.SourceEventID { + t.Fatalf("actual recorded media/identity changed: asset=%v proof=%v err=%v", stub.asset, stub.proof, err) + } + var result struct { + Outcome string `json:"outcome"` + Transcript []struct { + Role string `json:"role"` + Text string `json:"text"` + } `json:"transcript"` + Recording struct { + Status string `json:"status"` + } `json:"recording"` + } + if err := json.Unmarshal(stub.result, &result); err != nil || result.Outcome != "answered" || len(result.Transcript) != 1 || result.Transcript[0].Role != "user" || result.Transcript[0].Text != "synthetic final ASR fixture" || result.Recording.Status != "uploaded" { + t.Fatalf("current result does not reflect actual Mock media: %+v err=%v", result, err) + } + if entries, err := os.ReadDir(root); err != nil || len(entries) != 0 { + t.Fatalf("successful direct upload wrote business files: entries=%d err=%v", len(entries), err) + } +} + +func TestApprovedRecordedMockCallGenerationFailureReportsEmptyRecording(t *testing.T) { + runner, stub, puts, _, approved := recordedMockFixture(t, nil, 200*time.Millisecond) + err := runner.Run(context.Background(), approved) + if err == nil { + t.Fatal("Mock call with no observed media reported healthy completion") + } + if strings.Join(stub.calls, ",") != "end,result" || puts.Load() != 0 || stub.proof != nil { + t.Fatalf("missing media fabricated an upload or lost the terminal result: calls=%v puts=%d", stub.calls, puts.Load()) + } + var result struct { + Outcome string `json:"outcome"` + ReasonMessage string `json:"reason_message"` + Recording map[string]any `json:"recording"` + } + if err := json.Unmarshal(stub.result, &result); err != nil || result.Outcome != "failed" || !strings.Contains(result.ReasonMessage, "recording generation failed") || len(result.Recording) != 0 { + t.Fatalf("no-media failure did not produce an honest unique result: %+v err=%v", result, err) + } +} + +func TestApprovedRecordedMockCallRefusesMissingPerCallDelivery(t *testing.T) { + runner, stub, puts, _, approved := recordedMockFixture(t, bytes.Repeat([]byte{1, 0}, 1600), time.Second) + runner.Delivery = nil + if err := runner.Run(context.Background(), approved); err == nil || !errors.Is(err, ErrApprovedMockDeliveryUnavailable) || len(stub.calls) != 0 || puts.Load() != 0 { + t.Fatalf("missing per-call delivery was silently accepted: calls=%v err=%v", stub.calls, err) + } +}