Add Agent recording-fact unary client
This commit is contained in:
@@ -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 通过。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user