From 9b1f53953804d535b8f969e979fe6cbaf6c4235f Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 04:05:03 +0800 Subject: [PATCH] Verify task control over pinned mutual TLS --- .../saas-dispatcher-implementation.md | 2 +- .../rpc/approved_control_transport_test.go | 109 ++++++++++++++++++ 2 files changed, 110 insertions(+), 1 deletion(-) create mode 100644 internal/rpc/approved_control_transport_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 449a891..105c2cc 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -58,7 +58,7 @@ - Agent 隔离媒体入口 `rpc.RunApprovedCall` 仅接收签发的执行快照、媒体会话、挂断动作和显式每通话 AI pipeline;缺失或 typed nil 适配器在读媒体前拒绝,无隐式 SDK/Mock 回退。测试分别由已绑定快照构建 SDK pipeline 与注入受控合成脚本的 `ApprovedMockPipeline`。签发通话期限与 AI 总期限均约束整通电话;ASR-only 不合成开场,full-AI 开场完成才采集语音。父级期限不会被“未检测到语音”掩盖,空媒体 Mock 等到上下文真正结束。当前媒体只处理 16 kHz PCM16,SDK 虽接受、但媒体不能正确处理的 24 kHz ASR/TTS 配置在 Dispatcher 与 Agent 绑定时明确拒绝,不静默播放或识别错速音频。该入口尚未连接主 CLI,也未用于真实拨号。 - 隔离 Mock AI 适配器仅消费显式冻结的模拟媒体与最终识别脚本,不手写或替代供应商协议;缺失/多余轮次、未批准的开场或助理音频明确拒绝。ASR-only 不执行开场与回复;full-AI 开场只允许一次,用户最终识别的已配置字面关键词先挂断、绝不播放该轮候选回复,未配置关键词不推断拒联。任务未配置最大轮数时严格使用共享通话流程的一轮边界,负数仍拒绝。Mock 输出不代表真实 ASR/LLM/TTS 已通过。 - 新增 Agent 任务级在途通话栅栏 `TaskCalls`:按 D/数字租户/任务隔离;暂停或停止先拒新通话,hangup 请求取消、drain 不挂断,两者均等待已登记通话结束才确认;等待超时仍保留关闭栅栏,停止后不能恢复,同一控制重复执行不另设去重。隔离测试验证两任务互不干扰、重复结束不破坏状态。该栅栏已由活跃会话校验的 Agent RPC 隔离单测调用,但尚未接入实际批准通话的登记/取消和主入口;Dispatcher 的持久控制仍为权威,不能据此称真实通话已受控。 -- 新增内部任务级 `ApplyApprovedTaskControl` Proto、生成物及 Dispatcher `ApprovedOriginator.SendControl`:D 每次发送前从已激活会话领取元数据,只生成独立 RPC 追踪 ID,不从 SaaS 控制虚构 command_id、幂等键或 revision;Agent 校验活跃 D 会话、数字租户、任务和 pause/stop 的 hangup/drain 政策,在 Mock 中等待已登记活动通话结束后才确认,超时保留本地停止栅栏。隔离测试覆盖拒绝伪 D/旧会话、非法字段与政策、无适配器、未知 RPC 结果不自动重试及重复控制的显式重送;**尚无实际批准通话登记或新主 CLI 接线,也未经过该控制 RPC 的双向 TLS 传输测试**。 +- 新增内部任务级 `ApplyApprovedTaskControl` Proto、生成物及 Dispatcher `ApprovedOriginator.SendControl`:D 每次发送前从已激活会话领取元数据,只生成独立 RPC 追踪 ID,不从 SaaS 控制虚构 command_id、幂等键或 revision;Agent 校验活跃 D 会话、数字租户、任务和 pause/stop 的 hangup/drain 政策,在 Mock 中等待已登记活动通话结束后才确认,超时保留本地停止栅栏。隔离测试覆盖拒绝伪 D/旧会话、非法字段与政策、无适配器、未知 RPC 结果不自动重试及重复控制的显式重送;本地双向 TLS gRPC 还验证了受信 D 证书经真实生成 Stub 传送控制后,先挂断再等已登记通话结束才确认,不受信的客户端证书与旧会话均拒绝且不能重新准入。**尚无实际批准通话登记或新主 CLI 接线**,不能据此称真实通话控制已验收。 - 验证:`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`、`go test ./... -count=1`、`go test -race ./internal/ai ./internal/callflow ./internal/configread ./internal/dispatcher ./internal/rpc -count=1`、`go vet ./...`、`go build ./...`、`git diff --check` 均通过。 - **后续边界:** 录音直传/OSS 失败恢复及最终结果属 P06;新主 CLI 接线、真实媒体完整联动及旧路径清理属 P07;隔离端到端和 A01–A12 属 P08。真实 Agent/Asterisk、SaaS、MQ、AI 供应商联调未开展,不能由 Mock 结果代签。 diff --git a/internal/rpc/approved_control_transport_test.go b/internal/rpc/approved_control_transport_test.go new file mode 100644 index 0000000..6fa0674 --- /dev/null +++ b/internal/rpc/approved_control_transport_test.go @@ -0,0 +1,109 @@ +package rpc + +import ( + "context" + "errors" + "net" + "path/filepath" + "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/dispatcher" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/credentials" + "google.golang.org/grpc/status" + "google.golang.org/grpc/test/bufconn" +) + +func TestApprovedTaskControlOverMutualTLSRequiresPinnedPeerAndSession(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + execution := approvedTestRequest(t, now) + server := activatedApprovedServer(t, now, filepath.Join(t.TempDir(), "agent-session.json"), execution.DispatcherId, 1, nil) + calls := &agent.TaskCalls{} + server.approvedTaskCalls = calls + server.requirePeerCertificate = true + caPEM, caCert, caKey := testCertificate(t, nil, nil, true, nil, nil) + agentPEM, _, _ := testCertificate(t, caCert, caKey, false, []string{"agent.local"}, nil) + dispatcherPEM, dispatcherCert, _ := testCertificate(t, caCert, caKey, false, []string{"dispatcher.local"}, nil) + server.peerCertificateFingerprints = map[string]struct{}{CertificateFingerprint(dispatcherCert): {}} + serverTLS, err := NewServerTLSConfig(caPEM.certPEM, agentPEM.certPEM, agentPEM.keyPEM) + if err != nil { + t.Fatal(err) + } + 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() }) + dial := func(certPEM, keyPEM []byte) agentpb.AgentControlServiceClient { + t.Helper() + clientTLS, err := NewClientTLSConfig(caPEM.certPEM, certPEM, keyPEM, "agent.local") + if err != nil { + t.Fatal(err) + } + 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() }) + return agentpb.NewAgentControlServiceClient(conn) + } + client := dial(dispatcherPEM.certPEM, dispatcherPEM.keyPEM) + originator := &dispatcher.ApprovedOriginator{DispatcherID: execution.DispatcherId, Client: client, + Meta: func(context.Context) (*agentpb.RequestMeta, error) { return testMeta("task-control", "", 1), nil }, + } + task := agent.TaskIdentity{DispatcherID: execution.DispatcherId, TenantID: execution.TenantId, TaskID: execution.TaskId} + callCtx, cancelCall := context.WithCancel(context.Background()) + defer cancelCall() + release, err := calls.Register(task, "running-call", cancelCall) + if err != nil { + t.Fatal(err) + } + defer release() + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + spec := dispatcher.CurrentControlSpec{DispatcherID: task.DispatcherID, TenantID: task.TenantID, TaskID: task.TaskID, Action: "pause", ActiveCallPolicy: "hangup"} + completed := make(chan error, 1) + go func() { completed <- originator.SendControl(ctx, spec) }() + select { + case <-callCtx.Done(): + case <-ctx.Done(): + t.Fatal("authenticated control did not request hangup") + } + if _, err := calls.Register(task, "another-call", func() {}); !errors.Is(err, agent.ErrTaskAdmissionClosed) { + t.Fatalf("authenticated control did not fence new Agent work: %v", err) + } + select { + case err := <-completed: + t.Fatalf("Agent control returned before the physical call ended: %v", err) + default: + } + release() + if err := <-completed; err != nil { + t.Fatalf("mTLS control did not confirm the released call: %v", err) + } + foreignPEM, _, _ := testCertificate(t, caCert, caKey, false, []string{"other-dispatcher.local"}, nil) + foreignClient := dial(foreignPEM.certPEM, foreignPEM.keyPEM) + request := &agentpb.ApplyApprovedTaskControlRequest{ + Meta: testMeta("foreign-control", "", 1), DispatcherId: task.DispatcherID, + TenantId: task.TenantID, TaskId: task.TaskID, + Action: agentpb.ControlAction_CONTROL_ACTION_RESUME, + ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_UNSPECIFIED, + } + if _, err := foreignClient.ApplyApprovedTaskControl(ctx, request); status.Code(err) != codes.PermissionDenied { + t.Fatalf("unapproved client certificate was allowed to resume: %v", err) + } + request.Meta.SessionGeneration = 2 + if _, err := client.ApplyApprovedTaskControl(ctx, request); status.Code(err) != codes.Aborted { + t.Fatalf("stale authenticated session resumed a paused task: %v", err) + } + if _, err := calls.Register(task, "late-call", func() {}); !errors.Is(err, agent.ErrTaskAdmissionClosed) { + t.Fatalf("rejected mTLS controls changed Agent admission state: %v", err) + } +}