From 79c701dfd19a5c4910ccaf5c95a6edb609ca3b53 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 01:16:32 +0800 Subject: [PATCH] Add bounded in-memory WAV capture for callflow --- .../saas-dispatcher-implementation.md | 5 +- internal/callflow/recording.go | 81 ++++++++++++++++ internal/callflow/recording_test.go | 94 +++++++++++++++++++ internal/media/wav.go | 44 +++++++++ internal/media/wav_test.go | 38 ++++++++ 5 files changed, 260 insertions(+), 2 deletions(-) create mode 100644 internal/callflow/recording.go create mode 100644 internal/callflow/recording_test.go create mode 100644 internal/media/wav.go create mode 100644 internal/media/wav_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index aa985a8..2958ffb 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -64,8 +64,9 @@ - Dispatcher 的无录音最终结果隔离组件:`CurrentStore.RecordCallResult` 仅在确认通话结束后,按持久任务快照校验任务、被叫、主叫和已选线路,并以源执行事件固定生成唯一最终结果身份;消息通过严格 MQ Schema 校验后与 outbox 在同一事务写入。同内容重投/重启只恢复原消息,冲突结果和 SQLite 写入失败均不会产生第二份结果。这里只验证隔离组件,Agent 实际回报尚未连通。 - Dispatcher 的原始 OSS 目标及已上传结果隔离组件:新 SQLite 布局把一次通话的 upload_id、bucket、object_key、录音格式/时长/大小和 SHA-256 唯一绑定到已保留的执行;不保存临时 URL 或 TOKEN。旧布局拒绝启动并原样保留待交付 outbox,不自动迁移或清理。已签发录音目标不能通过空录音结果绕过上传;已有空录音最终结果不能再签发录音目标。`RecordUploadedCallResult` 仅接受与持久绑定完全一致的录音事实及 Agent 所报告的成功 PUT 状态,录音确认与唯一最终结果 outbox 同事务提交;丢失回报或 MQ 投递时重用原消息,已确认后拒绝再次签发 PUT 授权。Mock 证明的是本地状态约束,不是独立 OSS 校验或真实 Agent 身份验证。 - Dispatcher 的终结与外呼回执顺序竞争隔离修复:原流程在 Agent 接受执行的 RPC 返回后才写入外呼回执,快速结束或 RPC 超时可能先到;现以一次 SQLite 事务在确认通话已结束后补齐原回执并释放占用,未知执行仅在确认结束后释放。迟到的执行响应、超时和重复结束不会产生第二份回执;注入 outbox 写入失败保留原占用。并发竞争及结束后立即生成唯一最终结果有单元测试;主入口真实 Agent 会话注入与通话执行仍未接线。 -- 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、媒体录音或主入口接线。 -- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent ./internal/rpc ./internal/store -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↔Dispatcher 实际会话与录音执行接线(隔离 mTLS gRPC 已测)及上传事实/最终结果交付、MQ/端到端验收,不能宣称 P06 通过。 +- 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 通话的全量媒体验收。 +- 已验证:`go test ./... -count=1`、`go test -race ./internal/media ./internal/callflow -count=1`、`go test -race ./internal/agent ./internal/rpc ./internal/store -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↔Dispatcher 实际会话与录音执行接线(隔离 mTLS gRPC 已测)及上传事实/最终结果交付、MQ/端到端验收,不能宣称 P06 通过。 ## 验收台账 diff --git a/internal/callflow/recording.go b/internal/callflow/recording.go new file mode 100644 index 0000000..fcdbbdc --- /dev/null +++ b/internal/callflow/recording.go @@ -0,0 +1,81 @@ +package callflow + +import ( + "context" + "errors" + "fmt" + "sync" + + "git.ipao.vip/rogee/go-sip/internal/media" +) + +// RecordingSession records only PCM16 frames the approved callflow actually +// reads or successfully sends. It never writes a business file or encodes RTP. +// Capture failure is retained and returned by WAV; the call can still finish. +type RecordingSession struct { + MediaSession + mu sync.Mutex + pcm []byte + maxBytes int64 + failure error +} + +func NewRecordingSession(session MediaSession, maxWAVBytes int64) (*RecordingSession, error) { + if session == nil || maxWAVBytes <= 44 { + return nil, errors.New("recording requires a media session and a positive bounded WAV allowance") + } + return &RecordingSession{MediaSession: session, maxBytes: maxWAVBytes}, nil +} + +func (r *RecordingSession) ReadPayload(ctx context.Context) ([]byte, error) { + payload, err := r.MediaSession.ReadPayload(ctx) + if err == nil { + r.capture(payload) + } + return payload, err +} + +func (r *RecordingSession) SendPCM16(ctx context.Context, pcm []byte, sampleRateHz int) error { + if sampleRateHz != 16000 { + return fmt.Errorf("recording requires 16-kHz PCM16, received %d Hz", sampleRateHz) + } + if err := r.MediaSession.SendPCM16(ctx, pcm, sampleRateHz); err != nil { + return err + } + r.capture(pcm) + return nil +} + +func (r *RecordingSession) capture(pcm []byte) { + if len(pcm) == 0 { + return + } + r.mu.Lock() + defer r.mu.Unlock() + if r.failure != nil { + return + } + if len(pcm)%2 != 0 { + r.failure = media.ErrWAVInvalidPCM + r.pcm = nil + return + } + if int64(len(pcm)) > r.maxBytes-44-int64(len(r.pcm)) { + r.failure = media.ErrWAVTooLarge + r.pcm = nil + return + } + r.pcm = append(r.pcm, pcm...) +} + +// WAV returns an in-memory 16-kHz mono recording and its duration. The caller +// must inspect its error after the call; a failed capture must not be reported +// as a successful upload or a healthy no-recording outcome. +func (r *RecordingSession) WAV() ([]byte, int64, error) { + r.mu.Lock() + defer r.mu.Unlock() + if r.failure != nil { + return nil, 0, r.failure + } + return media.EncodeMonoWAV(r.pcm, r.maxBytes) +} diff --git a/internal/callflow/recording_test.go b/internal/callflow/recording_test.go new file mode 100644 index 0000000..eee1372 --- /dev/null +++ b/internal/callflow/recording_test.go @@ -0,0 +1,94 @@ +package callflow + +import ( + "bytes" + "context" + "errors" + "os" + "testing" + + "git.ipao.vip/rogee/go-sip/internal/media" +) + +func TestRecordingSessionCapturesReadAndSentMediaInMemory(t *testing.T) { + inbound := bytes.Repeat([]byte{1, 0}, 320) + outbound := bytes.Repeat([]byte{2, 0}, 320) + root := t.TempDir() + source := NewMemorySession(inbound) + recording, err := NewRecordingSession(source, 44+int64(len(inbound)+len(outbound))) + if err != nil { + t.Fatal(err) + } + if err := recording.SendPCM16(context.Background(), outbound, 16000); err != nil { + t.Fatal(err) + } + got, err := recording.ReadPayload(context.Background()) + if err != nil || !bytes.Equal(got, inbound) { + t.Fatalf("recording intercepted or changed the live inbound audio: err=%v", err) + } + wav, durationMS, err := recording.WAV() + if err != nil || durationMS != 40 || !bytes.Equal(wav[44:], append(bytes.Clone(outbound), inbound...)) { + t.Fatalf("recording did not preserve media in callflow order: duration=%d err=%v", durationMS, err) + } + if stats := recording.Stats(); stats.ReceivedPackets != 1 || stats.SentPackets != 1 { + t.Fatalf("recording changed live RTP accounting: %+v", stats) + } + entries, err := os.ReadDir(root) + if err != nil || len(entries) != 0 { + t.Fatalf("normal in-memory recording created business files: entries=%d err=%v", len(entries), err) + } +} + +func TestRecordingSessionReportsNoAudioAndCaptureFailuresWithoutInventingWAV(t *testing.T) { + if _, err := NewRecordingSession(nil, 1024); err == nil { + t.Fatal("missing media session admitted") + } + if _, err := NewRecordingSession(NewMemorySession(nil), 44); err == nil { + t.Fatal("no room for even one audio sample admitted") + } + empty, err := NewRecordingSession(NewMemorySession(nil), 1024) + if err != nil { + t.Fatal(err) + } + if _, _, err := empty.WAV(); !errors.Is(err, media.ErrWAVEmpty) { + t.Fatalf("no call audio became a recording: %v", err) + } + tooLarge, err := NewRecordingSession(NewMemorySession([]byte{3, 0, 4, 0}), 48) + if err != nil { + t.Fatal(err) + } + if err := tooLarge.SendPCM16(context.Background(), []byte{1, 0, 2, 0}, 16000); err != nil { + t.Fatal(err) + } + if frame, err := tooLarge.ReadPayload(context.Background()); err != nil || len(frame) != 4 { + t.Fatalf("recording size failure blocked the live call: frame=%v err=%v", frame, err) + } + if _, _, err := tooLarge.WAV(); !errors.Is(err, media.ErrWAVTooLarge) { + t.Fatalf("over-limit call audio became an uploadable WAV: %v", err) + } + odd, err := NewRecordingSession(NewMemorySession([]byte{1, 2, 3}), 1024) + if err != nil { + t.Fatal(err) + } + if frame, err := odd.ReadPayload(context.Background()); err != nil || len(frame) != 3 { + t.Fatalf("invalid recorded PCM blocked the original media stream: frame=%v err=%v", frame, err) + } + if _, _, err := odd.WAV(); !errors.Is(err, media.ErrWAVInvalidPCM) { + t.Fatalf("odd PCM16 was reported as a valid recording: %v", err) + } +} + +func TestRecordingSessionDoesNotRecordUnsentAudio(t *testing.T) { + recording, err := NewRecordingSession(NewMemorySession(nil), 1024) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := recording.SendPCM16(ctx, []byte{1, 0}, 16000); !errors.Is(err, context.Canceled) { + t.Fatalf("canceled media send unexpectedly succeeded: %v", err) + } + if _, _, err := recording.WAV(); !errors.Is(err, media.ErrWAVEmpty) { + t.Fatalf("unsent audio appeared in recording: %v", err) + } +} diff --git a/internal/media/wav.go b/internal/media/wav.go new file mode 100644 index 0000000..302285c --- /dev/null +++ b/internal/media/wav.go @@ -0,0 +1,44 @@ +package media + +import ( + "encoding/binary" + "errors" +) + +var ErrWAVEmpty = errors.New("no captured PCM16 recording") +var ErrWAVInvalidPCM = errors.New("captured PCM16 has an odd byte length") +var ErrWAVTooLarge = errors.New("captured WAV exceeds its configured byte limit") + +// EncodeMonoWAV wraps already decoded mono PCM16/16-kHz audio in memory. +// It is only a WAV container writer: RTP decoding and audio codecs remain in +// the selected media libraries. No business file is created on this path. +func EncodeMonoWAV(pcm []byte, maxBytes int64) ([]byte, int64, error) { + const headerBytes = 44 + if len(pcm) == 0 { + return nil, 0, ErrWAVEmpty + } + if len(pcm)%2 != 0 { + return nil, 0, ErrWAVInvalidPCM + } + if maxBytes <= headerBytes || int64(len(pcm)) > maxBytes-headerBytes || uint64(len(pcm)) > uint64(^uint32(0))-36 { + return nil, 0, ErrWAVTooLarge + } + const sampleRate = 16000 + dataSize := uint32(len(pcm)) + wav := make([]byte, headerBytes+len(pcm)) + copy(wav[:4], "RIFF") + binary.LittleEndian.PutUint32(wav[4:8], 36+dataSize) + copy(wav[8:12], "WAVE") + copy(wav[12:16], "fmt ") + binary.LittleEndian.PutUint32(wav[16:20], 16) + binary.LittleEndian.PutUint16(wav[20:22], 1) + binary.LittleEndian.PutUint16(wav[22:24], 1) + binary.LittleEndian.PutUint32(wav[24:28], sampleRate) + binary.LittleEndian.PutUint32(wav[28:32], sampleRate*2) + binary.LittleEndian.PutUint16(wav[32:34], 2) + binary.LittleEndian.PutUint16(wav[34:36], 16) + copy(wav[36:40], "data") + binary.LittleEndian.PutUint32(wav[40:44], dataSize) + copy(wav[44:], pcm) + return wav, int64(len(pcm)) * 1000 / (sampleRate * 2), nil +} diff --git a/internal/media/wav_test.go b/internal/media/wav_test.go new file mode 100644 index 0000000..26b526a --- /dev/null +++ b/internal/media/wav_test.go @@ -0,0 +1,38 @@ +package media + +import ( + "bytes" + "encoding/binary" + "testing" +) + +func TestEncodeMonoWAVProducesBoundedInMemoryPCM16(t *testing.T) { + pcm := bytes.Repeat([]byte{0x12, 0x34}, 3200) // 200 ms at 16 kHz. + wav, durationMS, err := EncodeMonoWAV(pcm, int64(len(pcm)+44)) + if err != nil { + t.Fatal(err) + } + if durationMS != 200 || len(wav) != 44+len(pcm) || string(wav[:4]) != "RIFF" || string(wav[8:12]) != "WAVE" || string(wav[12:16]) != "fmt " || string(wav[36:40]) != "data" { + t.Fatalf("in-memory recording does not contain the original PCM16: duration=%d bytes=%d", durationMS, len(wav)) + } + if binary.LittleEndian.Uint32(wav[4:8]) != uint32(len(wav)-8) || binary.LittleEndian.Uint16(wav[20:22]) != 1 || binary.LittleEndian.Uint16(wav[22:24]) != 1 || binary.LittleEndian.Uint32(wav[24:28]) != 16000 || binary.LittleEndian.Uint32(wav[28:32]) != 32000 || binary.LittleEndian.Uint16(wav[34:36]) != 16 || binary.LittleEndian.Uint32(wav[40:44]) != uint32(len(pcm)) || !bytes.Equal(wav[44:], pcm) { + t.Fatal("WAV header or copied PCM differs from the 16-kHz mono original") + } +} + +func TestEncodeMonoWAVRejectsEmptyOddAndOverLimitAudio(t *testing.T) { + for _, tc := range []struct { + name string + pcm []byte + limit int64 + }{ + {"no recording", nil, 1024}, + {"odd PCM16", []byte{1, 2, 3}, 1024}, + {"no storage allowance", []byte{1, 2}, 44}, + {"audio over bound", []byte{1, 2, 3, 4}, 47}, + } { + if _, _, err := EncodeMonoWAV(tc.pcm, tc.limit); err == nil { + t.Fatalf("%s was fabricated as an uploadable WAV", tc.name) + } + } +}