From 7bcde1812fb3fd952c3211b4dfd2d0d94ca5b932 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 07:21:16 +0800 Subject: [PATCH] Advance durable Agent session generations on restart --- .../saas-dispatcher-implementation.md | 2 + internal/dispatcher/agent.go | 7 ++- internal/dispatcher/agent_test.go | 8 +++- internal/rpc/server.go | 9 ++++ internal/rpc/session_test.go | 43 +++++++++++++++++++ 5 files changed, 66 insertions(+), 3 deletions(-) diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 1191b64..63d144d 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -36,6 +36,8 @@ - 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 进程重启或真实部署证据。 + ## P03:HTTP 读取分批改造(未整体签收) - `contract.ValidateCurrent` 与 `configread` 按当前 Schema 读取 SIP、provider、task、quota 和 cursor 任务发现;严格检查数字 tenant_id、本 D 归属及不可变配置。provider 凭据原值只保留在内存快照,不写日志;Agent 参数中的显式 0/false 保真;无旧 Schema/旧配置回退。 diff --git a/internal/dispatcher/agent.go b/internal/dispatcher/agent.go index 49ca390..facbeeb 100644 --- a/internal/dispatcher/agent.go +++ b/internal/dispatcher/agent.go @@ -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" + "github.com/google/uuid" "google.golang.org/protobuf/proto" ) @@ -156,7 +157,11 @@ func (c *AgentCoordinator) Activate(ctx context.Context, agentID, cellID, bootID if cellID == "" || bootID == "" || epoch == "" { return AgentSession{}, errors.New("cell, boot and dispatcher epoch are required") } - operationID := fmt.Sprintf("activate:%s:%s", agentID, bootID) + activationID, err := uuid.NewRandom() + if err != nil { + return AgentSession{}, fmt.Errorf("create Agent activation identity: %w", err) + } + 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}) if err != nil { diff --git a/internal/dispatcher/agent_test.go b/internal/dispatcher/agent_test.go index e1eff83..368bdd1 100644 --- a/internal/dispatcher/agent_test.go +++ b/internal/dispatcher/agent_test.go @@ -55,8 +55,12 @@ func TestAgentCoordinatorProbesBeforeActivation(t *testing.T) { if err != nil { t.Fatal(err) } - if first.SessionGeneration != 1 || second.SessionGeneration != 2 { - t.Fatalf("unexpected generations: first=%d second=%d", first.SessionGeneration, second.SessionGeneration) + third, err := coordinator.Activate(context.Background(), "agent-1", "cell-1", status.BootId, "epoch-2", 0) + if err != nil { + t.Fatal(err) + } + if first.SessionGeneration != 1 || second.SessionGeneration != 2 || third.SessionGeneration != 3 { + t.Fatalf("a new activation replayed an old session: first=%d second=%d third=%d", first.SessionGeneration, second.SessionGeneration, third.SessionGeneration) } } diff --git a/internal/rpc/server.go b/internal/rpc/server.go index 27ab969..42e89d8 100644 --- a/internal/rpc/server.go +++ b/internal/rpc/server.go @@ -294,6 +294,9 @@ func (r *SessionRegistry) Activate(binding *agentpb.AgentBinding, activationOper return cloneSession(existing.session), true, nil } if binding.SessionGeneration == 0 { + if existing.binding.SessionGeneration == ^uint64(0) { + return nil, false, status.Error(codes.Aborted, "session generation is exhausted") + } binding = proto.Clone(binding).(*agentpb.AgentBinding) binding.SessionGeneration = existing.binding.SessionGeneration + 1 } @@ -304,6 +307,12 @@ func (r *SessionRegistry) Activate(binding *agentpb.AgentBinding, activationOper if binding.SessionGeneration == 0 { binding = proto.Clone(binding).(*agentpb.AgentBinding) binding.SessionGeneration = 1 + if previous, ok := r.generations[binding.AgentId]; ok { + if previous == ^uint64(0) { + return nil, false, status.Error(codes.Aborted, "session generation is exhausted") + } + binding.SessionGeneration = previous + 1 + } } if previous, ok := r.generations[binding.AgentId]; ok && binding.SessionGeneration <= previous { return nil, false, status.Error(codes.Aborted, "persisted session generation is fenced") diff --git a/internal/rpc/session_test.go b/internal/rpc/session_test.go index 5ca3227..835580f 100644 --- a/internal/rpc/session_test.go +++ b/internal/rpc/session_test.go @@ -1,6 +1,7 @@ package rpc import ( + "strings" "testing" "time" @@ -25,3 +26,45 @@ func TestSessionRegistryPersistsGenerationAcrossRestart(t *testing.T) { t.Fatal(err) } } + +func TestSessionRegistryZeroGenerationAdvancesDurableHighwaterAfterRestart(t *testing.T) { + path := t.TempDir() + "/rpc-session.json" + first := NewSessionRegistry(path) + now := time.Unix(100, 0) + if _, _, err := first.Activate(&agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 7}, "activate-1", "digest-1", now); err != nil { + t.Fatal(err) + } + second := NewSessionRegistry(path) + session, replay, err := second.Activate(&agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-2", DispatcherEpoch: "epoch-2"}, "activate-2", "digest-2", now) + if err != nil || replay || session.GetSessionGeneration() != 8 { + t.Fatalf("Agent could not advance persisted fencing after restart: session=%v replay=%t err=%v", session, replay, err) + } + third := NewSessionRegistry(path) + if _, _, err := third.Activate(&agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-3", DispatcherEpoch: "epoch-3", SessionGeneration: 7}, "activate-3", "digest-3", now); status.Code(err) != codes.Aborted { + t.Fatalf("explicit stale generation bypassed the durable fence: %v", err) + } +} + +func TestSessionRegistryZeroGenerationDoesNotWrapAtExhaustion(t *testing.T) { + path := t.TempDir() + "/rpc-session.json" + first := NewSessionRegistry(path) + now := time.Unix(100, 0) + if _, _, err := first.Activate(&agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: ^uint64(0)}, "activate-1", "digest-1", now); err != nil { + t.Fatal(err) + } + for _, tc := range []struct { + name string + registry *SessionRegistry + boot string + }{ + {"active", first, "boot-1"}, + {"after restart", NewSessionRegistry(path), "boot-2"}, + } { + t.Run(tc.name, func(t *testing.T) { + _, _, err := tc.registry.Activate(&agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: tc.boot, DispatcherEpoch: "epoch-2"}, "activate-2", "digest-2", now) + if status.Code(err) != codes.Aborted || !strings.Contains(err.Error(), "exhausted") { + t.Fatalf("exhausted generation wrapped or was hidden: %v", err) + } + }) + } +}