diff --git a/cmd/sip-go-agent/current_agent_command_test.go b/cmd/sip-go-agent/current_agent_command_test.go index c65c3d4..85dc4c2 100644 --- a/cmd/sip-go-agent/current_agent_command_test.go +++ b/cmd/sip-go-agent/current_agent_command_test.go @@ -101,6 +101,7 @@ func TestCurrentAgentCommandServesPinnedMutualTLSOnly(t *testing.T) { "AGENT_ID": settings.AgentID, "CELL_ID": settings.CellID, "AGENT_GRPC_LISTEN": settings.Listen, "AGENT_SESSION_PATH": settings.SessionPath, "AGENT_RECOVERY_ROOT": settings.RecoveryRoot, + "DISPATCHER_ID": settings.DispatcherID, "DISPATCHER_GRPC_ENDPOINT": settings.DispatcherEndpoint, "DISPATCHER_GRPC_SERVER_NAME": settings.DispatcherServerName, "MTLS_CA_FILE": settings.CAFile, "MTLS_CERT_FILE": settings.CertFile, @@ -118,8 +119,12 @@ func TestCurrentAgentCommandServesPinnedMutualTLSOnly(t *testing.T) { command.SetArgs([]string{"--mode", "mock"}) finished := make(chan error, 1) go func() { finished <- command.Execute() }() + finishedConsumed := false t.Cleanup(func() { cancel() + if finishedConsumed { + return + } select { case err := <-finished: if err != nil { @@ -139,6 +144,7 @@ func TestCurrentAgentCommandServesPinnedMutualTLSOnly(t *testing.T) { if err != nil { select { case startupErr := <-finished: + finishedConsumed = true t.Fatalf("Agent startup failed: %v (dial: %v)", startupErr, err) default: t.Fatalf("isolated Agent mTLS listener unavailable: %v", err) @@ -149,10 +155,30 @@ func TestCurrentAgentCommandServesPinnedMutualTLSOnly(t *testing.T) { Meta: &agentpb.RequestMeta{ProtocolVersion: "agent.v1", RequestId: "mock-status", TraceId: "mock-status", OperationId: "mock-status", AgentId: settings.AgentID, CellId: settings.CellID}, Target: &agentpb.AgentBinding{AgentId: settings.AgentID, CellId: settings.CellID}, } - response, err := agentpb.NewAgentControlServiceClient(trusted).GetAgentStatus(dial, probe) + trustedClient := agentpb.NewAgentControlServiceClient(trusted) + response, err := trustedClient.GetAgentStatus(dial, probe) if err != nil || response.GetStatus().GetAgentId() != settings.AgentID || response.GetStatus().GetCellId() != settings.CellID { t.Fatalf("trusted D could not probe the actual Agent command: response=%v err=%v", response, err) } + activation := &agentpb.ActivateAgentRequest{ + Meta: &agentpb.RequestMeta{ProtocolVersion: "agent.v1", RequestId: "bind-dispatcher", TraceId: "bind-dispatcher", + OperationId: "bind-dispatcher", AgentId: settings.AgentID, CellId: settings.CellID, + BootId: response.GetStatus().GetBootId(), DispatcherEpoch: "epoch-test"}, + Binding: &agentpb.AgentBinding{AgentId: settings.AgentID, CellId: settings.CellID, + ExpectedBootId: response.GetStatus().GetBootId(), DispatcherEpoch: "epoch-test", + DispatcherId: "22222222-2222-4222-8222-222222222222"}, + ActivationOperationId: "bind-dispatcher", + } + if _, err := trustedClient.ActivateAgent(dial, activation); status.Code(err) != codes.PermissionDenied { + t.Fatalf("pinned certificate could claim another Dispatcher identity: %v", err) + } + if _, err := os.Stat(settings.SessionPath); !os.IsNotExist(err) { + t.Fatalf("unapproved Dispatcher wrote a session: %v", err) + } + activation.Binding.DispatcherId = settings.DispatcherID + if _, err := trustedClient.ActivateAgent(dial, activation); err != nil { + t.Fatalf("approved Dispatcher could not activate the Agent: %v", err) + } untrustedTLS, err := rpc.NewClientTLSConfig(ca, agentCert, agentKey, "agent.local") if err != nil { t.Fatal(err) diff --git a/cmd/sip-go-agent/current_agent_setup.go b/cmd/sip-go-agent/current_agent_setup.go index 6b74588..4d12e1a 100644 --- a/cmd/sip-go-agent/current_agent_setup.go +++ b/cmd/sip-go-agent/current_agent_setup.go @@ -18,6 +18,7 @@ import ( "git.ipao.vip/rogee/go-sip/internal/agent" "git.ipao.vip/rogee/go-sip/internal/config" "git.ipao.vip/rogee/go-sip/internal/rpc" + "git.ipao.vip/rogee/go-sip/internal/tenant" "github.com/google/uuid" ) @@ -30,6 +31,9 @@ func newCurrentAgentServer(ctx context.Context, settings config.AgentEnvironment scenario.MaxWAVBytes <= 44 || len(scenario.Script.Turns) == 0 || strings.TrimSpace(scenario.ReasonMessage) == "" { return nil, errors.New("current Agent requires an active process, explicit Mock media and deployment identity") } + if tenant.ValidateDispatcherID(settings.DispatcherID) != nil { + return nil, errors.New("current Agent requires an approved Dispatcher UUID v4") + } if dispatcher == nil { return nil, errors.New("current Agent requires a pinned Dispatcher transport") } @@ -92,7 +96,7 @@ func newCurrentAgentServer(ctx context.Context, settings config.AgentEnvironment }, } handler, err = rpc.NewApprovedAgentServer(rpc.ServerOptions{ - Mode: "mock", StatePath: settings.SessionPath, + Mode: "mock", StatePath: settings.SessionPath, ApprovedDispatcherID: settings.DispatcherID, Status: &agentpb.AgentStatus{AgentId: settings.AgentID, CellId: settings.CellID, BootId: uuid.NewString(), ProtocolVersion: "agent.v1"}, LoadedSIP: func(context.Context) (map[string]int64, error) { return maps.Clone(loaded), nil }, PeerCertificateFingerprints: pins, RequirePeerCertificate: true, diff --git a/cmd/sip-go-agent/current_agent_setup_test.go b/cmd/sip-go-agent/current_agent_setup_test.go index 16b8361..49f8edf 100644 --- a/cmd/sip-go-agent/current_agent_setup_test.go +++ b/cmd/sip-go-agent/current_agent_setup_test.go @@ -26,7 +26,7 @@ func currentAgentSetupFixture(t *testing.T) (config.AgentEnvironment, approvedMo t.Fatal(err) } settings := config.AgentEnvironment{ - AgentID: "agent-mock", CellID: "cell-mock", Listen: "127.0.0.1:0", + AgentID: "agent-mock", CellID: "cell-mock", DispatcherID: "11111111-1111-4111-8111-111111111111", Listen: "127.0.0.1:0", SessionPath: filepath.Join(root, "session.json"), RecoveryRoot: root, DispatcherEndpoint: "127.0.0.1:39443", DispatcherServerName: "dispatcher.local", CAFile: filepath.Join(root, "ca.pem"), CertFile: filepath.Join(root, "agent.pem"), @@ -60,6 +60,9 @@ func TestNewCurrentAgentServerRefusesUnsafeAdaptersBeforeResources(t *testing.T) {"missing Agent identity", func(c *config.AgentEnvironment, _ *approvedMockScenario, _ *map[string]int64, _ *agentpb.AgentControlServiceClient) { c.AgentID = "" }}, + {"missing authorized Dispatcher", func(c *config.AgentEnvironment, _ *approvedMockScenario, _ *map[string]int64, _ *agentpb.AgentControlServiceClient) { + c.DispatcherID = "" + }}, {"missing SIP evidence", func(_ *config.AgentEnvironment, _ *approvedMockScenario, sip *map[string]int64, _ *agentpb.AgentControlServiceClient) { *sip = nil }}, diff --git a/deploys/env/agent.env.example b/deploys/env/agent.env.example index 2651523..8b8c578 100644 --- a/deploys/env/agent.env.example +++ b/deploys/env/agent.env.example @@ -4,6 +4,8 @@ # Start explicitly: sip-go-agent agent --mode mock AGENT_ID=agent-mock CELL_ID=cell-mock +# Required: the preapproved canonical UUID v4 of this Agent's sole Dispatcher. +# DISPATCHER_ID= AGENT_GRPC_LISTEN=127.0.0.1:19090 AGENT_SESSION_PATH=/var/lib/sip-go-agent/agent/session.json # Pre-create this directory with mode 0700; neither state nor files are cleared. diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 889e407..1191b64 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -34,7 +34,7 @@ - 新增无实现代次字段的 Dispatcher 环境预检:必须显式给出规范 UUID v4 身份、只读 HTTP 地址及密钥、RabbitMQ 地址和 SQLite 路径;当前 Mock 仅接受本机 HTTP/MQ 目标,mixed/real 在任何资源操作前拒绝。预检不打开数据库或网络,缺失值不继承旧默认配置,错误不打印凭据。`go test ./internal/config -run '^TestLoadDispatcherEnvironment' -count=1` 通过;此预检尚未接入主 CLI,不能当作 P02/P03 完成。 -- Agent 增加独立的无代次环境预检:单 Agent/Cell 身份、会话与私有恢复路径、隔离 Mock 场景和已加载 SIP 测试事实、D gRPC 目标及 mTLS 文件/指纹均须显式提供;仅允许本机监听和本机 D 端点,mixed/real 在访问文件或网络前拒绝。缺失项不继承旧 `FromEnv` 默认值,凭据/地址不回显;`go test ./internal/config -run '^TestLoadAgentEnvironment' -count=1` 通过。此检查已接入根命令可达的唯一 Agent Mock 入口;本机双向 TLS 监听可由受信 D 证书探测,另一张同 CA 证书被指纹门禁拒绝。该测试尚未激活 D 会话、执行呼叫、验证 Agent→D 录音或真实 Asterisk 加载。 +- Agent 增加独立的无代次环境预检:Agent/Cell 与预授权 D 的身份、会话与私有恢复路径、隔离 Mock 场景和已加载 SIP 测试事实、D gRPC 目标及 mTLS 文件/指纹均须显式提供;仅允许本机监听和本机 D 端点,mixed/real 在访问文件或网络前拒绝。缺失项不继承旧 `FromEnv` 默认值,凭据/地址不回显;`go test ./internal/config -run '^TestLoadAgentEnvironment' -count=1` 通过。此检查已接入根命令可达的唯一 Agent Mock 入口;本机双向 TLS 监听可由受信 D 证书探测,另一张同 CA 证书被指纹门禁拒绝。本地测试用批准的 D 身份激活会话,错误 D UUID 即使携带受信证书也在写入会话日志前拒绝;尚未执行呼叫、验证 Agent→D 录音或真实 Asterisk 加载。 ## P03:HTTP 读取分批改造(未整体签收) diff --git a/internal/config/agent_runtime.go b/internal/config/agent_runtime.go index 4694d75..f420274 100644 --- a/internal/config/agent_runtime.go +++ b/internal/config/agent_runtime.go @@ -7,6 +7,8 @@ import ( "os" "strconv" "strings" + + "git.ipao.vip/rogee/go-sip/internal/tenant" ) // AgentEnvironment is deployment-owned and is inspected without opening any @@ -14,6 +16,7 @@ import ( type AgentEnvironment struct { AgentID string CellID string + DispatcherID string Listen string SessionPath string RecoveryRoot string @@ -44,6 +47,7 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { value *string }{ {"AGENT_ID", &settings.AgentID}, {"CELL_ID", &settings.CellID}, + {"DISPATCHER_ID", &settings.DispatcherID}, {"AGENT_GRPC_LISTEN", &settings.Listen}, {"AGENT_SESSION_PATH", &settings.SessionPath}, {"AGENT_RECOVERY_ROOT", &settings.RecoveryRoot}, {"DISPATCHER_GRPC_ENDPOINT", &settings.DispatcherEndpoint}, @@ -59,6 +63,9 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { } *field.value = value } + if tenant.ValidateDispatcherID(settings.DispatcherID) != nil { + return AgentEnvironment{}, errors.New("DISPATCHER_ID must be a canonical UUID v4") + } if !localGRPCAddress(settings.Listen, true) { return AgentEnvironment{}, errors.New("AGENT_GRPC_LISTEN must be an isolated local Mock address") } diff --git a/internal/config/agent_runtime_test.go b/internal/config/agent_runtime_test.go index a9fa782..a39a11f 100644 --- a/internal/config/agent_runtime_test.go +++ b/internal/config/agent_runtime_test.go @@ -14,6 +14,7 @@ func setAgentEnvironment(t *testing.T) string { for name, value := range map[string]string{ "AGENT_ID": "agent-mock", "CELL_ID": "cell-mock", "AGENT_GRPC_LISTEN": "127.0.0.1:0", "AGENT_SESSION_PATH": state, "AGENT_RECOVERY_ROOT": filepath.Join(root, "recovery"), + "DISPATCHER_ID": "11111111-1111-4111-8111-111111111111", "DISPATCHER_GRPC_ENDPOINT": "127.0.0.1:39443", "DISPATCHER_GRPC_SERVER_NAME": "dispatcher.local", "MTLS_CA_FILE": filepath.Join(root, "ca.pem"), "MTLS_CERT_FILE": filepath.Join(root, "agent.pem"), "MTLS_KEY_FILE": filepath.Join(root, "agent.key"), "MTLS_PEER_CERT_FINGERPRINTS": strings.Repeat("a", 64), @@ -42,7 +43,7 @@ func TestLoadAgentEnvironmentRefusesNonMockBeforeResources(t *testing.T) { func TestLoadAgentEnvironmentRequiresExplicitDeploymentValues(t *testing.T) { for _, name := range []string{ "AGENT_ID", "CELL_ID", "AGENT_GRPC_LISTEN", "AGENT_SESSION_PATH", "AGENT_RECOVERY_ROOT", - "DISPATCHER_GRPC_ENDPOINT", "DISPATCHER_GRPC_SERVER_NAME", "MTLS_CA_FILE", "MTLS_CERT_FILE", + "DISPATCHER_ID", "DISPATCHER_GRPC_ENDPOINT", "DISPATCHER_GRPC_SERVER_NAME", "MTLS_CA_FILE", "MTLS_CERT_FILE", "MTLS_KEY_FILE", "MTLS_PEER_CERT_FINGERPRINTS", "AGENT_MOCK_SCENARIO_FILE", "AGENT_MOCK_APPLIED_SIP_FILE", } { t.Run(name, func(t *testing.T) { @@ -61,7 +62,7 @@ func TestLoadAgentEnvironmentRequiresLocalMockEndpointsAndDoesNotOpenFiles(t *te if err != nil { t.Fatal(err) } - if settings.AgentID != "agent-mock" || settings.CellID != "cell-mock" || settings.Listen != "127.0.0.1:0" || settings.SessionPath != state || settings.DispatcherEndpoint != "127.0.0.1:39443" { + if settings.AgentID != "agent-mock" || settings.CellID != "cell-mock" || settings.DispatcherID != "11111111-1111-4111-8111-111111111111" || settings.Listen != "127.0.0.1:0" || settings.SessionPath != state || settings.DispatcherEndpoint != "127.0.0.1:39443" { t.Fatal("explicit deployment identity or address was rewritten") } if _, err := os.Stat(state); !os.IsNotExist(err) { @@ -69,6 +70,7 @@ func TestLoadAgentEnvironmentRequiresLocalMockEndpointsAndDoesNotOpenFiles(t *te } for _, tc := range []struct{ name, value string }{ {"AGENT_GRPC_LISTEN", "0.0.0.0:39444"}, + {"DISPATCHER_ID", "unapproved-dispatcher"}, {"DISPATCHER_GRPC_ENDPOINT", "saas.example.invalid:443"}, } { t.Run(tc.name, func(t *testing.T) { diff --git a/internal/rpc/approved_server.go b/internal/rpc/approved_server.go index 491a648..0b201cd 100644 --- a/internal/rpc/approved_server.go +++ b/internal/rpc/approved_server.go @@ -1,6 +1,9 @@ package rpc -import "errors" +import ( + "errors" + "strings" +) var ErrApprovedWorkerRequired = errors.New("approved call worker is required") @@ -12,8 +15,8 @@ func NewApprovedAgentServer(options ServerOptions, worker *ApprovedCallWorker) ( if worker == nil { return nil, ErrApprovedWorkerRequired } - if options.Mode != "mock" || options.StatePath == "" || options.Status == nil || options.Status.AgentId == "" || options.Status.CellId == "" || options.Status.BootId == "" || options.LoadedSIP == nil { - return nil, errors.New("approved Agent requires explicit Mock mode, durable state, identity and observed SIP") + if options.Mode != "mock" || options.StatePath == "" || options.Status == nil || options.Status.AgentId == "" || options.Status.CellId == "" || options.Status.BootId == "" || strings.TrimSpace(options.ApprovedDispatcherID) == "" || options.LoadedSIP == nil { + return nil, errors.New("approved Agent requires explicit Mock mode, durable state, Agent and Dispatcher identity, and observed SIP") } if worker.Lifecycle == nil || worker.Calls == nil || worker.Prepare == nil || worker.OnFailure == nil { return nil, errors.New("approved Agent requires a process lifecycle, task calls, runner and failure reporting") diff --git a/internal/rpc/approved_server_test.go b/internal/rpc/approved_server_test.go index 4aa14a3..6420f95 100644 --- a/internal/rpc/approved_server_test.go +++ b/internal/rpc/approved_server_test.go @@ -3,6 +3,7 @@ package rpc import ( "context" "errors" + "os" "path/filepath" "testing" "time" @@ -40,7 +41,7 @@ func TestApprovedAgentServerRegistersCallsBeforeExecuteAckAndDrainsOnControl(t * } server, err := NewApprovedAgentServer(ServerOptions{ Mode: "mock", StatePath: filepath.Join(t.TempDir(), "agent-session.json"), - Now: func() time.Time { return now }, + ApprovedDispatcherID: req.DispatcherId, Now: func() time.Time { return now }, Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, LoadedSIP: func(context.Context) (map[string]int64, error) { return map[string]int64{req.SelectedTrunkId: req.SipRevision}, nil @@ -111,7 +112,7 @@ func TestApprovedAgentServerRefusesMissingOrCompetingAdapters(t *testing.T) { OnFailure: func(ApprovedExecution, error) error { return nil }, } valid := ServerOptions{Mode: "mock", StatePath: filepath.Join(t.TempDir(), "agent-session.json"), - Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, + ApprovedDispatcherID: "dispatcher-1", Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, LoadedSIP: func(context.Context) (map[string]int64, error) { return map[string]int64{"trunk-mock": 8}, nil }, } for _, tc := range []struct { @@ -123,6 +124,7 @@ func TestApprovedAgentServerRefusesMissingOrCompetingAdapters(t *testing.T) { {"implicit Mock default", func(o *ServerOptions, _ *ApprovedCallWorker) { o.Mode = "" }}, {"missing SIP observation", func(o *ServerOptions, _ *ApprovedCallWorker) { o.LoadedSIP = nil }}, {"missing Agent identity", func(o *ServerOptions, _ *ApprovedCallWorker) { o.Status = nil }}, + {"missing authorized Dispatcher", func(o *ServerOptions, _ *ApprovedCallWorker) { o.ApprovedDispatcherID = "" }}, {"competing originator", func(o *ServerOptions, _ *ApprovedCallWorker) { o.MockApprovedOriginate = func(context.Context, ApprovedExecution) error { return nil } }}, @@ -143,6 +145,43 @@ func TestApprovedAgentServerRefusesMissingOrCompetingAdapters(t *testing.T) { } } +func TestApprovedAgentServerRejectsSelfReportedDispatcherBeforeSessionWrite(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + req := approvedTestRequest(t, now) + state := filepath.Join(t.TempDir(), "agent-session.json") + worker := &ApprovedCallWorker{Lifecycle: context.Background(), Calls: &agent.TaskCalls{}, + Prepare: prepareWorkerRun(func(context.Context, ApprovedExecution) error { return nil }), + OnFailure: func(ApprovedExecution, error) error { return nil }, + } + server, err := NewApprovedAgentServer(ServerOptions{ + Mode: "mock", StatePath: state, ApprovedDispatcherID: req.DispatcherId, + Now: func() time.Time { return now }, + Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, + LoadedSIP: func(context.Context) (map[string]int64, error) { + return map[string]int64{req.SelectedTrunkId: req.SipRevision}, nil + }, + }, worker) + if err != nil { + t.Fatal(err) + } + activation := &agentpb.ActivateAgentRequest{ + Meta: testMeta("activate-unapproved-d", "", 0), + Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", + DispatcherEpoch: "epoch-1", DispatcherId: req.DispatcherId + "-other"}, + ActivationOperationId: "activate-unapproved-d", + } + if _, err := server.ActivateAgent(context.Background(), activation); status.Code(err) != codes.PermissionDenied { + t.Fatalf("pinned peer could claim a different Dispatcher identity: %v", err) + } + if _, err := os.Stat(state); !os.IsNotExist(err) { + t.Fatalf("rejected identity persisted a session: %v", err) + } + activation.Binding.DispatcherId = req.DispatcherId + if _, err := server.ActivateAgent(context.Background(), activation); err != nil { + t.Fatalf("approved Dispatcher could not activate: %v", err) + } +} + func TestExecuteApprovedReturnsDefinitiveNoDialRefusalWithoutRetry(t *testing.T) { now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) for _, tc := range []struct { diff --git a/internal/rpc/server.go b/internal/rpc/server.go index a2b6871..27ab969 100644 --- a/internal/rpc/server.go +++ b/internal/rpc/server.go @@ -40,6 +40,7 @@ type ServerOptions struct { AIAuthorizationRaw []byte Now func() time.Time RequirePeerCertificate bool + ApprovedDispatcherID string PeerAgentIDs map[string]string PeerCertificateFingerprints map[string]struct{} StatePath string @@ -73,6 +74,7 @@ type Server struct { aiAuthorizationRaw []byte aiConfigError error requirePeerCertificate bool + approvedDispatcherID string peerAgentIDs map[string]string peerCertificateFingerprints map[string]struct{} callLogger *calllog.Logger @@ -176,6 +178,7 @@ func NewServer(options ServerOptions) *Server { aiAuthorizationRaw: append([]byte(nil), options.AIAuthorizationRaw...), aiConfigError: aiConfigError, requirePeerCertificate: options.RequirePeerCertificate, + approvedDispatcherID: options.ApprovedDispatcherID, peerAgentIDs: cloneStringMap(options.PeerAgentIDs), peerCertificateFingerprints: cloneSet(options.PeerCertificateFingerprints), callLogger: options.CallLogger, @@ -449,6 +452,9 @@ func (s *Server) ActivateAgent(ctx context.Context, req *agentpb.ActivateAgentRe if err := s.checkPeer(ctx, req.Binding.AgentId); err != nil { return nil, err } + if s.approvedDispatcherID != "" && req.Binding.DispatcherId != s.approvedDispatcherID { + return nil, status.Error(codes.PermissionDenied, "activated Dispatcher identity is not authorized") + } if s.staticArtifactEnabled { if s.staticArtifactError != nil { return nil, status.Errorf(codes.FailedPrecondition, "static Cell artifact is invalid: %v", s.staticArtifactError)