Verify task control over pinned mutual TLS
This commit is contained in:
@@ -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 结果代签。
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user