Require Dispatcher identity in Agent activation
This commit is contained in:
@@ -435,7 +435,7 @@ func connectConfiguredAgents(ctx context.Context, cfg config.Config) (*connected
|
||||
}
|
||||
// Zero lets Agent-side durable session state allocate the next generation
|
||||
// after a Dispatcher restart; hard-coding 1 would self-fence recovery.
|
||||
session, activateErr := runtime.coordinator.Activate(ctx, endpoint.AgentID, endpoint.CellID, status.BootId, epoch, 0)
|
||||
session, activateErr := runtime.coordinator.Activate(ctx, cfg.DispatcherID, endpoint.AgentID, endpoint.CellID, status.BootId, epoch, 0)
|
||||
if activateErr != nil {
|
||||
return fail(fmt.Errorf("activate Agent %q: %w", endpoint.AgentID, activateErr))
|
||||
}
|
||||
|
||||
@@ -36,6 +36,7 @@
|
||||
- Dispatcher 另有纯运行环境预检:唯一 Agent 端点清单和 OSS 配置文件须提供路径,本机 gRPC 监听及双向 TLS 文件/Agent 证书指纹须显式给出;缺失、非法指纹或非本机监听明确拒绝且不回显值。该检查不读取配置文件、不连接 Agent、不打开 SQLite/MQ;`go test ./internal/config -run '^TestLoadDispatcherRuntimeEnvironment' -count=1` 通过。当前仍未由 `dispatcher` 主命令调用,不能当作 Agent 加载或 OSS 签发已验收。
|
||||
- D 的部署文件读取已独立于旧代次 Schema:Agent 清单仅接受**一个本机端点**;OSS 配置必须匹配本 D 身份、指定本机 HTTPS 目标、15 分钟授权与显式资产上限,凭据只经受控环境变量引用。超大文件、重复/未知字段、多个 Agent、非本机目标、缺失凭据和错误 D 归属均拒绝;JSON 唯一键及凭据引用基础函数已移出旧配置文件。`go test ./internal/config -count=1` 通过。此处仅验证部署数据,不代表 Agent 实际加载、官方 SDK 已签发目标,亦未接入主 `dispatcher` 命令。
|
||||
|
||||
- D 的 Agent 会话激活现在必须显式携带规范 UUID v4 的本 D 身份,并写入 Agent 授权绑定;预配置了对应 D 身份的隔离 Agent 会拒绝错误归属,不能再靠 D 自报空身份放行。`go test ./internal/dispatcher -count=1` 通过;主 `dispatcher` 命令尚未使用这条会话链路,不代表本机 D↔A 执行已验收。
|
||||
- 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 加载。
|
||||
|
||||
- 会话栅栏修复:Agent 重启后 `session_generation=0` 由持久化代际高水位递增,旧显式代际仍拒绝,代际耗尽明确阻断;D 每次新激活产生独立操作身份,同 boot/epoch 的下一次激活不再误回放旧会话。隔离重启/重复激活测试通过,未作为主 Dispatcher 进程重启或真实部署证据。
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
"git.ipao.vip/rogee/go-sip/internal/store"
|
||||
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
||||
"github.com/google/uuid"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
@@ -149,7 +150,10 @@ func (c *AgentCoordinator) Probe(ctx context.Context, agentID, cellID string) (*
|
||||
return response.Status, nil
|
||||
}
|
||||
|
||||
func (c *AgentCoordinator) Activate(ctx context.Context, agentID, cellID, bootID, epoch string, generation uint64) (AgentSession, error) {
|
||||
func (c *AgentCoordinator) Activate(ctx context.Context, dispatcherID, agentID, cellID, bootID, epoch string, generation uint64) (AgentSession, error) {
|
||||
if tenant.ValidateDispatcherID(dispatcherID) != nil {
|
||||
return AgentSession{}, errors.New("Agent activation requires a canonical Dispatcher UUID v4")
|
||||
}
|
||||
client, err := c.client(agentID)
|
||||
if err != nil {
|
||||
return AgentSession{}, err
|
||||
@@ -163,11 +167,11 @@ func (c *AgentCoordinator) Activate(ctx context.Context, agentID, cellID, bootID
|
||||
}
|
||||
operationID := "activate:" + activationID.String()
|
||||
meta := &agentpb.RequestMeta{ProtocolVersion: "agent.v1", RequestId: operationID + ":request", TraceId: operationID, OperationId: operationID, DispatcherEpoch: epoch, AgentId: agentID, CellId: cellID, BootId: bootID}
|
||||
response, err := client.ActivateAgent(ctx, &agentpb.ActivateAgentRequest{Meta: meta, Binding: &agentpb.AgentBinding{AgentId: agentID, CellId: cellID, ExpectedBootId: bootID, DispatcherEpoch: epoch, SessionGeneration: generation}, ActivationOperationId: operationID})
|
||||
response, err := client.ActivateAgent(ctx, &agentpb.ActivateAgentRequest{Meta: meta, Binding: &agentpb.AgentBinding{AgentId: agentID, CellId: cellID, ExpectedBootId: bootID, DispatcherId: dispatcherID, DispatcherEpoch: epoch, SessionGeneration: generation}, ActivationOperationId: operationID})
|
||||
if err != nil {
|
||||
return AgentSession{}, err
|
||||
}
|
||||
if response.Session == nil || response.State != agentpb.ActivationState_ACTIVATION_STATE_ACTIVE {
|
||||
if response == nil || response.Session == nil || response.State != agentpb.ActivationState_ACTIVATION_STATE_ACTIVE {
|
||||
return AgentSession{}, errors.New("Agent activation was not active")
|
||||
}
|
||||
session := AgentSession{AgentID: agentID, CellID: cellID, BootID: bootID, DispatcherEpoch: epoch, SessionGeneration: response.Session.SessionGeneration, ExpiresAtUnixMs: response.Session.ExpiresAtUnixMs}
|
||||
|
||||
@@ -15,10 +15,13 @@ import (
|
||||
"google.golang.org/grpc/test/bufconn"
|
||||
)
|
||||
|
||||
const approvedTestDispatcherID = "c046b893-8628-4589-ae50-619d049248a6"
|
||||
|
||||
func TestAgentCoordinatorProbesBeforeActivation(t *testing.T) {
|
||||
listener := bufconn.Listen(1024 * 1024)
|
||||
server := rpcserver.NewServer(rpcserver.ServerOptions{
|
||||
Now: func() time.Time { return time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC) },
|
||||
ApprovedDispatcherID: approvedTestDispatcherID,
|
||||
Now: func() time.Time { return time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC) },
|
||||
Status: &agentpb.AgentStatus{
|
||||
AgentId: "agent-1",
|
||||
CellId: "cell-1",
|
||||
@@ -47,15 +50,18 @@ func TestAgentCoordinatorProbesBeforeActivation(t *testing.T) {
|
||||
if status.BootId != "boot-probed" || status.SessionActive {
|
||||
t.Fatalf("unexpected pre-activation status: %+v", status)
|
||||
}
|
||||
first, err := coordinator.Activate(context.Background(), "agent-1", "cell-1", status.BootId, "epoch-1", 0)
|
||||
if _, err := coordinator.Activate(context.Background(), "not-a-dispatcher", "agent-1", "cell-1", status.BootId, "epoch-1", 0); err == nil {
|
||||
t.Fatal("unapproved Dispatcher identity was admitted")
|
||||
}
|
||||
first, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-1", "cell-1", status.BootId, "epoch-1", 0)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
second, err := coordinator.Activate(context.Background(), "agent-1", "cell-1", status.BootId, "epoch-2", 0)
|
||||
second, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-1", "cell-1", status.BootId, "epoch-2", 0)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
third, err := coordinator.Activate(context.Background(), "agent-1", "cell-1", status.BootId, "epoch-2", 0)
|
||||
third, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-1", "cell-1", status.BootId, "epoch-2", 0)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -81,7 +87,7 @@ func TestAgentCoordinatorActivatesExecutesAndControlsWithoutRetryingOriginate(t
|
||||
if err := coordinator.Register("agent-1", agentpb.NewAgentControlServiceClient(conn)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := coordinator.Activate(context.Background(), "agent-1", "cell-1", "boot-1", "epoch-1", 1); err != nil {
|
||||
if _, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-1", "cell-1", "boot-1", "epoch-1", 1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
|
||||
@@ -266,7 +266,7 @@ func TestExecuteReservedBindsQuotaAndAgentExecution(t *testing.T) {
|
||||
if err := coordinator.Register("agent-1", client); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := coordinator.Activate(context.Background(), "agent-1", "cell-1", "boot-1", "epoch-1", 1); err != nil {
|
||||
if _, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-1", "cell-1", "boot-1", "epoch-1", 1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result, err := d.ExecuteReserved(context.Background(), coordinator, "agent-1", task, raw, "reservation-integration", "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa")
|
||||
|
||||
@@ -88,7 +88,7 @@ func TestLocalContractBackedFlowEvidence(t *testing.T) {
|
||||
if err := coordinator.Register("agent-a", client); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := coordinator.Activate(context.Background(), "agent-a", "cell-a", "boot-a", "epoch-a", 1); err != nil {
|
||||
if _, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-a", "cell-a", "boot-a", "epoch-a", 1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result, err := d.ExecuteReserved(context.Background(), coordinator, "agent-a", task, callRaw, "reservation-local-flow", snapshot.Digest)
|
||||
|
||||
@@ -41,7 +41,7 @@ func TestLocalAuthorizedOriginationReachesOnlyMockAgentAfterFinalGate(t *testing
|
||||
if err := coordinator.Register("agent-1", agentpb.NewAgentControlServiceClient(connection)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := coordinator.Activate(context.Background(), "agent-1", "cell-1", "boot-1", "epoch-1", 1); err != nil {
|
||||
if _, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-1", "cell-1", "boot-1", "epoch-1", 1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
executionID := acceptedLocalExecutionID(t, d, st, at, "command-agent-bridge")
|
||||
|
||||
@@ -26,10 +26,10 @@ func TestTwoMockCellsKeepAgentSessionsAndPermitsSeparate(t *testing.T) {
|
||||
if err := coordinator.Register("agent-b", clientB); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := coordinator.Activate(context.Background(), "agent-a", "cell-a", "boot-a", "epoch-1", 1); err != nil {
|
||||
if _, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-a", "cell-a", "boot-a", "epoch-1", 1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := coordinator.Activate(context.Background(), "agent-b", "cell-b", "boot-b", "epoch-1", 1); err != nil {
|
||||
if _, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-b", "cell-b", "boot-b", "epoch-1", 1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user