Bind Agent sessions to preapproved Dispatcher identity
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
}},
|
||||
|
||||
Vendored
+2
@@ -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=<injected-uuid-v4>
|
||||
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.
|
||||
|
||||
@@ -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 读取分批改造(未整体签收)
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user