Verify recording delivery across local TLS and SQLite
This commit is contained in:
@@ -64,10 +64,11 @@
|
||||
- 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 或主入口批准执行的媒体录音接线。
|
||||
- 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 边界约束。此处未连接真实 D gRPC、主入口批准执行媒体或 MQ。
|
||||
- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent ./internal/rpc ./internal/store ./internal/callflow ./internal/media -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 通过。
|
||||
- Agent 录音交付隔离组件:`RecordingDelivery` 先确认结束,再依照录音是否实际生成分别上报唯一空录音结果或请求原授权并直传内存 WAV;录音生成失败保留通话真实结果、空录音对象及明确原因,不虚构上传事实。隔离测试通过本地 HTTP PUT 和假 Dispatcher RPC 覆盖成功无业务文件、OSS 明确失败后私有文件保存、恢复写入失败、未知 PUT 隔离、重启重领原目标、上传已确认后只重发原结果。再次调用不会隐式重新 PUT;正常已确认上传但尚未被 D 持久收讫的跨进程间隙仍受 K16 边界约束。此处未连接主入口批准执行媒体或 MQ;下述隔离链路另测本地真正的 D gRPC/SQLite。
|
||||
- 本地隔离链路:`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 通过。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -0,0 +1,141 @@
|
||||
package rpc
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net"
|
||||
"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/callflow"
|
||||
"git.ipao.vip/rogee/go-sip/internal/oss"
|
||||
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials"
|
||||
"google.golang.org/grpc/test/bufconn"
|
||||
)
|
||||
|
||||
func TestRecordingDeliveryRealMutualTLSAndLocalOSSCommitsOneSQLiteResult(t *testing.T) {
|
||||
server, database, _, request, snapshot := recordingRPCFixture(t)
|
||||
mediaSession, err := callflow.NewRecordingSession(callflow.NewMemorySession(bytes.Repeat([]byte{1, 0}, 320)), 1024)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
frame, err := mediaSession.ReadPayload(context.Background())
|
||||
if err != nil || len(frame) != 640 {
|
||||
t.Fatalf("live media was not captured: size=%d err=%v", len(frame), err)
|
||||
}
|
||||
wav, durationMS, err := mediaSession.WAV()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var puts atomic.Int32
|
||||
localOSS := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
puts.Add(1)
|
||||
body, err := io.ReadAll(io.LimitReader(r.Body, 1025))
|
||||
if err != nil || r.Method != http.MethodPut || !strings.HasPrefix(r.URL.Path, "/mock-bucket/approved/") || !bytes.Equal(body, wav) {
|
||||
t.Errorf("local OSS received an incorrect PUT: method=%q size=%d path=%q err=%v", r.Method, len(body), r.URL.Path, err)
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusCreated)
|
||||
}))
|
||||
defer localOSS.Close()
|
||||
server.OSS, err = oss.NewClient(oss.Config{
|
||||
Endpoint: localOSS.URL, Region: "cn-test", Bucket: "mock-bucket", KeyPrefix: "approved",
|
||||
AccessKeyID: "isolated-test-key", AccessKeySecret: "isolated-test-secret",
|
||||
GrantTTL: 15 * time.Minute, MaxAssetBytes: 1024,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
caPEM, caCert, caKey := testCertificate(t, nil, nil, true, nil, nil)
|
||||
serverPEM, _, _ := testCertificate(t, caCert, caKey, false, []string{"dispatcher.local"}, nil)
|
||||
clientPEM, clientCert, _ := testCertificate(t, caCert, caKey, false, []string{"agent.local"}, nil)
|
||||
serverTLS, err := NewServerTLSConfig(caPEM.certPEM, serverPEM.certPEM, serverPEM.keyPEM)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
clientTLS, err := NewClientTLSConfig(caPEM.certPEM, clientPEM.certPEM, clientPEM.keyPEM, "dispatcher.local")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
server.TrustedFingerprints = map[string]struct{}{CertificateFingerprint(clientCert): {}}
|
||||
listener := bufconn.Listen(1 << 20)
|
||||
grpcServer := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS)))
|
||||
agentpb.RegisterAgentControlServiceServer(grpcServer, server)
|
||||
go func() { _ = grpcServer.Serve(listener) }()
|
||||
t.Cleanup(func() { grpcServer.Stop(); _ = listener.Close() })
|
||||
conn, err := grpc.NewClient("bufnet",
|
||||
grpc.WithContextDialer(func(context.Context, string) (net.Conn, error) { return listener.Dial() }),
|
||||
grpc.WithTransportCredentials(credentials.NewTLS(clientTLS)),
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = conn.Close() })
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||||
defer cancel()
|
||||
client := agent.RecordingClient{
|
||||
Client: agentpb.NewAgentControlServiceClient(conn), DispatcherID: request.DispatcherId,
|
||||
TenantID: request.TenantId, SourceEventID: request.SourceEventId,
|
||||
Session: func(context.Context) (*agentpb.RequestMeta, error) { return request.Meta, nil },
|
||||
}
|
||||
|
||||
var result map[string]any
|
||||
if err := json.Unmarshal(recordingResultPayload(t, snapshot, nil), &result); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result["outcome"] = "answered"
|
||||
result["reason_code"] = 200
|
||||
result["reason_message"] = "completed"
|
||||
payload, err := json.Marshal(result)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
root := t.TempDir()
|
||||
if err := os.Chmod(root, 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
delivery := &agent.RecordingDelivery{
|
||||
Call: client,
|
||||
Recovery: &agent.RecordingRecovery{
|
||||
Root: root, Now: server.Now,
|
||||
Upload: agent.UploadClient{HTTPClient: localOSS.Client(), Now: server.Now},
|
||||
},
|
||||
}
|
||||
if err := delivery.Complete(ctx, agent.CompletedRecording{
|
||||
ResultPayload: payload, Expected: true, RecordingID: request.Asset.AssetId,
|
||||
UploadID: request.UploadId, WAV: wav, DurationMS: durationMS,
|
||||
}); err != nil {
|
||||
t.Fatalf("local mutual-TLS Agent→D→OSS→SQLite recording flow failed: %v", err)
|
||||
}
|
||||
if puts.Load() != 1 {
|
||||
t.Fatalf("approved recording PUT occurred %d times", puts.Load())
|
||||
}
|
||||
stored, err := database.LoadRecordingUpload(request.DispatcherId, request.SourceEventId)
|
||||
if err != nil || stored.ConfirmedAt == "" || stored.SizeBytes != int64(len(wav)) || stored.Bucket != "mock-bucket" {
|
||||
t.Fatalf("the confirmed original upload and result did not commit together: confirmed=%t size=%d bucket=%q err=%v", stored.ConfirmedAt != "", stored.SizeBytes, stored.Bucket, err)
|
||||
}
|
||||
outbox, err := database.ListPendingOutbox(request.DispatcherId)
|
||||
if err != nil || len(outbox) != 2 || outbox[0].EventType != "call.execute" || outbox[1].EventType != "call.execute.result" {
|
||||
t.Fatalf("execution created multiple/missing outbound facts: count=%d err=%v", len(outbox), err)
|
||||
}
|
||||
entries, err := os.ReadDir(root)
|
||||
if err != nil || len(entries) != 0 {
|
||||
t.Fatalf("successful direct upload created local business files: files=%d err=%v", len(entries), err)
|
||||
}
|
||||
if request.Meta.OperationId != "request-1" || request.Meta.IdempotencyKey != "request-1" {
|
||||
t.Fatal("Agent session metadata was mutated by the recording flow")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user