From 1cc2d12b73a564533123a66a29afa162383f5cb6 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 01:03:18 +0800 Subject: [PATCH] Add Agent recording-fact unary client --- .../saas-dispatcher-implementation.md | 6 +- internal/agent/recording_client.go | 111 ++++++++++++++++++ internal/agent/recording_client_test.go | 21 ++++ internal/rpc/recording_transport_test.go | 69 +++++++++++ 4 files changed, 204 insertions(+), 3 deletions(-) create mode 100644 internal/agent/recording_client.go create mode 100644 internal/agent/recording_client_test.go create mode 100644 internal/rpc/recording_transport_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index ed32652..aa985a8 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -63,9 +63,9 @@ - Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;真实 Dispatcher 授权及回报尚未接线。 - 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、已上传与无录音结果及重复回报。当前是服务方法的本地隔离测试,尚无真实 gRPC 握手、Agent 调用或实际 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 受控会话/完整 gRPC 传输及上传事实/最终结果交付、MQ/端到端验收,不能宣称 P06 通过。 +- 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 通过。 ## 验收台账 diff --git a/internal/agent/recording_client.go b/internal/agent/recording_client.go new file mode 100644 index 0000000..24b7b33 --- /dev/null +++ b/internal/agent/recording_client.go @@ -0,0 +1,111 @@ +package agent + +import ( + "bytes" + "context" + "errors" + "fmt" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + + "google.golang.org/protobuf/proto" +) + +var ErrRecordingClientUnavailable = errors.New("Agent recording delivery requires an active Dispatcher session and Unary client") + +// RecordingClient belongs to one approved call. It requests upload tokens only +// when its caller explicitly asks, and it never retries an OSS PUT or invents a +// successful call end or result after an RPC error. +type RecordingClient struct { + Client agentpb.AgentControlServiceClient + Session func(context.Context) (*agentpb.RequestMeta, error) + DispatcherID string + TenantID int64 + SourceEventID string +} + +func (c RecordingClient) requestMeta(ctx context.Context, action string) (*agentpb.RequestMeta, error) { + if c.Client == nil || c.Session == nil || c.DispatcherID == "" || c.TenantID <= 0 || c.SourceEventID == "" { + return nil, ErrRecordingClientUnavailable + } + if err := ctx.Err(); err != nil { + return nil, err + } + meta, err := c.Session(ctx) + if err != nil { + return nil, fmt.Errorf("obtain current Agent session: %w", err) + } + if meta == nil || meta.GetAgentId() == "" || meta.GetCellId() == "" || meta.GetBootId() == "" || meta.GetDispatcherEpoch() == "" || meta.GetSessionGeneration() == 0 { + return nil, ErrRecordingClientUnavailable + } + copy := proto.Clone(meta).(*agentpb.RequestMeta) + copy.OperationId = c.SourceEventID + "/" + action + copy.IdempotencyKey = copy.OperationId + return copy, nil +} + +func (c RecordingClient) RequestUpload(ctx context.Context, asset *agentpb.AssetDescriptor, uploadID string) (*agentpb.UploadGrant, error) { + meta, err := c.requestMeta(ctx, "upload/"+uploadID) + if err != nil { + return nil, err + } + if asset == nil || uploadID == "" || asset.GetKind() != agentpb.AssetKind_ASSET_KIND_RECORDING || + asset.GetExecutionId() != c.SourceEventID || asset.GetCallId() != c.SourceEventID || asset.GetAssetId() == "" || asset.GetSizeBytes() <= 0 || asset.GetChecksumSha256() == "" { + return nil, errors.New("recording token request does not describe the approved call and original asset") + } + response, err := c.Client.RequestRecordingUpload(ctx, &agentpb.RequestRecordingUploadRequest{ + Meta: meta, DispatcherId: c.DispatcherID, TenantId: c.TenantID, SourceEventId: c.SourceEventID, + UploadId: uploadID, Asset: proto.Clone(asset).(*agentpb.AssetDescriptor), + }) + if err != nil { + return nil, fmt.Errorf("request original recording upload token: %w", err) + } + grant := response.GetGrant() + if grant == nil || grant.GetUploadId() != uploadID || grant.GetBucket() == "" || grant.GetObjectKey() == "" || + grant.GetMaxBytes() != asset.GetSizeBytes() || grant.GetRequiredChecksumSha256() != asset.GetChecksumSha256() { + return nil, fmt.Errorf("%w: Dispatcher returned a different recording asset", ErrUploadGrantInvalid) + } + return grant, nil +} + +func (c RecordingClient) ReportEnded(ctx context.Context) error { + meta, err := c.requestMeta(ctx, "ended") + if err != nil { + return err + } + response, err := c.Client.ReportCallEnded(ctx, &agentpb.ReportCallEndedRequest{ + Meta: meta, DispatcherId: c.DispatcherID, TenantId: c.TenantID, SourceEventId: c.SourceEventID, + }) + if err != nil { + return fmt.Errorf("report confirmed call end: %w", err) + } + if receipt := response.GetReceipt(); receipt.GetResult() != agentpb.ResultCode_RESULT_CODE_APPLIED || receipt.GetFactId() != c.SourceEventID { + return errors.New("Dispatcher did not persist the confirmed call end") + } + return nil +} + +func (c RecordingClient) ReportFinal(ctx context.Context, payload []byte, upload *agentpb.UploadObservation) (string, error) { + meta, err := c.requestMeta(ctx, "result") + if err != nil { + return "", err + } + if len(payload) == 0 { + return "", errors.New("final call result payload is required") + } + request := &agentpb.ReportCallResultRequest{ + Meta: meta, DispatcherId: c.DispatcherID, TenantId: c.TenantID, SourceEventId: c.SourceEventID, + ResultPayloadJson: bytes.Clone(payload), + } + if upload != nil { + request.Upload = proto.Clone(upload).(*agentpb.UploadObservation) + } + response, err := c.Client.ReportCallResult(ctx, request) + if err != nil { + return "", fmt.Errorf("persist unique call result: %w", err) + } + if receipt := response.GetReceipt(); receipt.GetResult() == agentpb.ResultCode_RESULT_CODE_ACCEPTED && receipt.GetFactId() != "" { + return receipt.GetFactId(), nil + } + return "", errors.New("Dispatcher did not persist the unique call result") +} diff --git a/internal/agent/recording_client_test.go b/internal/agent/recording_client_test.go new file mode 100644 index 0000000..7d2f6d0 --- /dev/null +++ b/internal/agent/recording_client_test.go @@ -0,0 +1,21 @@ +package agent + +import ( + "context" + "testing" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" +) + +func TestRecordingClientRejectsMissingAuthorizedDispatcherSession(t *testing.T) { + client := RecordingClient{} + if _, err := client.RequestUpload(context.Background(), &agentpb.AssetDescriptor{}, "upload-1"); err == nil { + t.Fatal("unconfigured Agent requested an OSS upload token") + } + if err := client.ReportEnded(context.Background()); err == nil { + t.Fatal("unconfigured Agent claimed a confirmed call end") + } + if _, err := client.ReportFinal(context.Background(), []byte(`{}`), nil); err == nil { + t.Fatal("unconfigured Agent sent a final result") + } +} diff --git a/internal/rpc/recording_transport_test.go b/internal/rpc/recording_transport_test.go new file mode 100644 index 0000000..3ee5b8b --- /dev/null +++ b/internal/rpc/recording_transport_test.go @@ -0,0 +1,69 @@ +package rpc + +import ( + "context" + "net" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/agent" + + "google.golang.org/grpc" + "google.golang.org/grpc/credentials" + "google.golang.org/grpc/test/bufconn" +) + +func TestRecordingFactsCrossRealMutualTLSUnaryTransportWithoutExternalOSS(t *testing.T) { + server, database, _, request, snapshot := recordingRPCFixture(t) + 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 }, + } + grant, err := client.RequestUpload(ctx, request.Asset, request.UploadId) + if err != nil || grant.GetBucket() != "mock-bucket" { + t.Fatalf("mTLS Agent request did not bind Dispatcher-owned bucket: bucket=%q err=%v", grant.GetBucket(), err) + } + if err := client.ReportEnded(ctx); err != nil { + t.Fatalf("mTLS Agent end did not durably acknowledge original execution: %v", err) + } + observation := &agentpb.UploadObservation{UploadId: request.UploadId, RecordingId: request.Asset.AssetId, PutStatusCode: 200, SizeBytes: request.Asset.SizeBytes, ChecksumSha256: request.Asset.ChecksumSha256} + resultID, err := client.ReportFinal(ctx, recordingResultPayload(t, snapshot, grant), observation) + if err != nil || resultID == "" { + t.Fatalf("mTLS Agent final result not persisted: fact=%q err=%v", resultID, err) + } + if request.Meta.OperationId != "request-1" || request.Meta.IdempotencyKey != "request-1" { + t.Fatal("concurrent Agent session metadata was mutated by the recording client") + } + 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("gRPC facts did not produce one ACK and one result: outbox=%+v err=%v", outbox, err) + } +}