fix(dispatcher): renew approved Agent session before expiry

This commit is contained in:
2026-09-30 18:18:37 +08:00
parent 9089cbc787
commit 66062b0c6c
6 changed files with 342 additions and 3 deletions
+23 -1
View File
@@ -116,7 +116,8 @@ func runDispatcher(ctx context.Context, mode string) (result error) {
cancelProbe()
return fmt.Errorf("create Dispatcher epoch: %w", err)
}
if _, err := coordinator.Activate(probeCtx, settings.DispatcherID, agentEndpoint.AgentID, agentEndpoint.CellID, status.BootId, epoch.String(), 0); err != nil {
session, err := coordinator.Activate(probeCtx, settings.DispatcherID, agentEndpoint.AgentID, agentEndpoint.CellID, status.BootId, epoch.String(), 0)
if err != nil {
cancelProbe()
return fmt.Errorf("activate assigned Agent before admission: %w", err)
}
@@ -152,14 +153,35 @@ func runDispatcher(ctx context.Context, mode string) (result error) {
})
serveCtx, cancelServe := context.WithCancel(ctx)
defer cancelServe()
renewErrors := make(chan error, 1)
go func() {
renewErr := maintainAgentSession(serveCtx, session, 5*time.Minute, func(renewCtx context.Context, previous dispatcher.AgentSession) (dispatcher.AgentSession, error) {
return coordinator.Renew(renewCtx, settings.DispatcherID, previous)
})
if renewErr != nil {
closeErr := database.CloseAdmission(settings.DispatcherID)
renewErrors <- errors.Join(renewErr, closeErr)
cancelServe()
return
}
renewErrors <- nil
}()
serverErrors := make(chan error, 1)
go func() {
serverErrors <- server.Serve(listener)
cancelServe()
}()
runtimeErr := worker.Serve(serveCtx)
cancelServe()
server.GracefulStop()
serverErr := <-serverErrors
renewErr := <-renewErrors
if renewErr != nil {
if errors.Is(serverErr, grpc.ErrServerStopped) {
serverErr = nil
}
return errors.Join(runtimeErr, fmt.Errorf("Agent session renewal failed: %w", renewErr), serverErr)
}
if serverErr != nil && !errors.Is(serverErr, grpc.ErrServerStopped) {
return errors.Join(runtimeErr, fmt.Errorf("Dispatcher recording listener stopped: %w", serverErr))
}
+60
View File
@@ -0,0 +1,60 @@
package main
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
)
// maintainAgentSession renews before the ten-minute Agent session expires.
// A failed or invalid renewal stops admission rather than using the old lease.
func maintainAgentSession(ctx context.Context, active dispatcher.AgentSession, before time.Duration, renew func(context.Context, dispatcher.AgentSession) (dispatcher.AgentSession, error)) error {
if ctx == nil || renew == nil || before <= 0 {
return errors.New("Agent session renewal requires a context, interval and renewal operation")
}
for {
if ctx.Err() != nil {
return nil
}
if active.ExpiresAtUnixMs <= time.Now().UnixMilli() {
return fmt.Errorf("Agent %s session generation %d expired before renewal", active.AgentID, active.SessionGeneration)
}
wait := time.Until(time.UnixMilli(active.ExpiresAtUnixMs).Add(-before))
if wait < 0 {
wait = 0
}
timer := time.NewTimer(wait)
select {
case <-ctx.Done():
timer.Stop()
return nil
case <-timer.C:
}
if ctx.Err() != nil {
return nil
}
if active.ExpiresAtUnixMs <= time.Now().UnixMilli() {
return fmt.Errorf("Agent %s session generation %d expired before renewal", active.AgentID, active.SessionGeneration)
}
requestCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
next, err := renew(requestCtx, active)
cancel()
if err != nil {
if ctx.Err() != nil {
return nil
}
return fmt.Errorf("renew Agent %s session generation %d: %w", active.AgentID, active.SessionGeneration, err)
}
if next.AgentID != active.AgentID || next.CellID != active.CellID || next.BootID != active.BootID ||
next.DispatcherEpoch != active.DispatcherEpoch || next.SessionGeneration <= active.SessionGeneration ||
next.ExpiresAtUnixMs <= active.ExpiresAtUnixMs || next.ExpiresAtUnixMs <= time.Now().Add(before).UnixMilli() {
return fmt.Errorf("Agent %s session generation %d renewal returned no usable session", active.AgentID, active.SessionGeneration)
}
slog.Info("Agent session renewed", "agent_id", next.AgentID, "session_generation", next.SessionGeneration, "expires_at_unix_ms", next.ExpiresAtUnixMs)
active = next
}
}
@@ -0,0 +1,85 @@
package main
import (
"context"
"errors"
"strings"
"testing"
"time"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
)
func testAgentSession(expiresAt time.Time) dispatcher.AgentSession {
return dispatcher.AgentSession{
AgentID: "agent-1", CellID: "cell-1", BootID: "boot-1", DispatcherEpoch: "epoch-1",
SessionGeneration: 1, ExpiresAtUnixMs: expiresAt.UnixMilli(),
}
}
func TestMaintainAgentSessionRenewsBeforeExpiryAndReportsFailure(t *testing.T) {
initial := testAgentSession(time.Now().Add(3 * time.Second))
failure := errors.New("Agent unavailable")
calls := 0
err := maintainAgentSession(context.Background(), initial, 2*time.Second, func(_ context.Context, previous dispatcher.AgentSession) (dispatcher.AgentSession, error) {
calls++
if time.Now().UnixMilli() >= previous.ExpiresAtUnixMs {
t.Fatal("session expired before renewal attempt")
}
if calls == 1 {
if previous != initial {
t.Fatalf("first renewal changed identity: %+v", previous)
}
next := previous
next.SessionGeneration++
next.ExpiresAtUnixMs = time.Now().Add(3 * time.Second).UnixMilli()
return next, nil
}
return dispatcher.AgentSession{}, failure
})
if calls != 2 || !errors.Is(err, failure) {
t.Fatalf("session renewal did not repeat and report a definite failure: calls=%d err=%v", calls, err)
}
}
func TestMaintainAgentSessionFailsClosedWithoutAUsableNextSession(t *testing.T) {
for name, next := range map[string]func(dispatcher.AgentSession) dispatcher.AgentSession{
"stale generation": func(previous dispatcher.AgentSession) dispatcher.AgentSession { return previous },
"changed boot": func(previous dispatcher.AgentSession) dispatcher.AgentSession {
previous.BootID = "unknown-boot"
previous.SessionGeneration++
previous.ExpiresAtUnixMs = time.Now().Add(time.Minute).UnixMilli()
return previous
},
} {
t.Run(name, func(t *testing.T) {
initial := testAgentSession(time.Now().Add(2 * time.Second))
calls := 0
err := maintainAgentSession(context.Background(), initial, time.Second, func(_ context.Context, previous dispatcher.AgentSession) (dispatcher.AgentSession, error) {
calls++
return next(previous), nil
})
if err == nil || calls != 1 {
t.Fatalf("invalid renewal advanced admission or looped: calls=%d err=%v", calls, err)
}
})
}
called := false
if err := maintainAgentSession(context.Background(), testAgentSession(time.Now().Add(-time.Second)), time.Second, func(context.Context, dispatcher.AgentSession) (dispatcher.AgentSession, error) {
called = true
return dispatcher.AgentSession{}, nil
}); err == nil || !strings.Contains(err.Error(), "expired") || called {
t.Fatalf("an expired session was silently renewed: err=%v called=%t", err, called)
}
}
func TestMaintainAgentSessionStopsOnShutdown(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
if err := maintainAgentSession(ctx, testAgentSession(time.Now().Add(10*time.Minute)), 5*time.Minute, func(context.Context, dispatcher.AgentSession) (dispatcher.AgentSession, error) {
t.Fatal("shutdown started another activation")
return dispatcher.AgentSession{}, nil
}); err != nil {
t.Fatalf("normal shutdown reported an activation failure: %v", err)
}
}
@@ -155,3 +155,7 @@
## 验收台账
P01–P07 的项目内隔离证据见上;P07 唯一当前入口、全仓残留/合法例外及来源/hash 的逐项结果见 [`saas-dispatcher-p07-audit.md`](saas-dispatcher-p07-audit.md)。A01–A12 和 K01–K16 的最终逐项对照、业务覆盖率及诊断门禁仍待 P08。不得用本地 Mock 冒充真实 SaaS、management、MQ、OSS、AI、Asterisk、SIP、ECS 或生产签收。
## P08 本地验收进度
- 十分钟 Agent 会话续期:审查发现 Dispatcher 之前只在启动时激活一次,约十分钟后 Agent 和 Dispatcher 均拒绝过期会话,长时间运行时无法再执行/上报。新增从真实会话到期时刻计算的提前五分钟续期;只允许同一个已审批 Agent boot、Cell 与 Dispatcher epoch 生成更高代际,激活响应、到期前/状态探测后的有效期及报告中的代际均须复核。续期失败立即关闭新准入、取消服务并以 Agent ID/代际的脱敏错误说明原因,不接纳未知新 boot、不自动重拨。TDD 先复现没有 `Renew`/循环、状态探测期间过期仍被重新激活,再以 bufconn 真实双端会话、可控时钟及重复 race 测试证明原会话过期后新会话仍可操作、旧代际被拒绝、换 boot/过期/失败不续;正常取消不再发起激活。`go test -race ./internal/dispatcher ./cmd/sip-go-agent -count=1`、`make check`(三项隔离 MQ 明确 PASS)、`make coverage`(全部非生成手写代码语句覆盖率 **71.8%**)、`make release-check-local` 通过。录音上报与续期切换竞争的故障证据和非生产诊断脚本门禁仍待审查,不将此条视为 P08 完成或真实主机验证。
+38 -2
View File
@@ -101,8 +101,9 @@ func (c *AgentCoordinator) Activate(ctx context.Context, dispatcherID, agentID,
if err != nil {
return AgentSession{}, err
}
if response == nil || response.Session == nil || response.State != agentpb.ActivationState_ACTIVATION_STATE_ACTIVE {
return AgentSession{}, errors.New("Agent activation was not active")
if response == nil || response.Session == nil || response.State != agentpb.ActivationState_ACTIVATION_STATE_ACTIVE ||
response.Session.DispatcherEpoch != epoch || response.Session.SessionGeneration == 0 || response.Session.ExpiresAtUnixMs <= c.now().UnixMilli() {
return AgentSession{}, errors.New("Agent activation returned an invalid or expired session")
}
session := AgentSession{AgentID: agentID, CellID: cellID, BootID: bootID, DispatcherEpoch: epoch, SessionGeneration: response.Session.SessionGeneration, ExpiresAtUnixMs: response.Session.ExpiresAtUnixMs}
c.mu.Lock()
@@ -111,6 +112,41 @@ func (c *AgentCoordinator) Activate(ctx context.Context, dispatcherID, agentID,
return session, nil
}
// Renew reactivates only the existing approved Agent boot before its session
// expires. Failure never adopts another boot or grants a second execution.
func (c *AgentCoordinator) Renew(ctx context.Context, dispatcherID string, previous AgentSession) (AgentSession, error) {
if c == nil || ctx == nil || previous.AgentID == "" || previous.CellID == "" || previous.BootID == "" || previous.DispatcherEpoch == "" {
return AgentSession{}, errors.New("complete approved Agent session is required for renewal")
}
if err := ctx.Err(); err != nil {
return AgentSession{}, err
}
c.mu.Lock()
active, ok := c.sessions[previous.AgentID]
c.mu.Unlock()
if !ok || active != previous || previous.ExpiresAtUnixMs <= c.now().UnixMilli() {
return AgentSession{}, errors.New("Agent session changed or expired before renewal")
}
status, err := c.Probe(ctx, previous.AgentID, previous.CellID)
if err != nil {
return AgentSession{}, fmt.Errorf("probe Agent before session renewal: %w", err)
}
if status.BootId != previous.BootID {
return AgentSession{}, errors.New("Agent boot changed during session renewal; unknown work requires disposition")
}
if previous.ExpiresAtUnixMs <= c.now().UnixMilli() {
return AgentSession{}, errors.New("Agent session expired during status probing")
}
renewed, err := c.Activate(ctx, dispatcherID, previous.AgentID, previous.CellID, previous.BootID, previous.DispatcherEpoch, 0)
if err != nil {
return AgentSession{}, fmt.Errorf("renew approved Agent session: %w", err)
}
if renewed.SessionGeneration <= previous.SessionGeneration || renewed.ExpiresAtUnixMs <= previous.ExpiresAtUnixMs {
return AgentSession{}, errors.New("Agent renewal did not advance generation and expiry")
}
return renewed, nil
}
func (c *AgentCoordinator) client(agentID string) (agentpb.AgentControlServiceClient, error) {
c.mu.Lock()
defer c.mu.Unlock()
+132
View File
@@ -0,0 +1,132 @@
package dispatcher
import (
"context"
"net"
"testing"
"time"
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
rpcserver "git.ipao.vip/rogee/go-sip/internal/rpc"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/test/bufconn"
)
func TestAgentCoordinatorRenewsTenMinuteSessionWithoutAdoptingAnotherBoot(t *testing.T) {
clock := time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC)
listener := bufconn.Listen(1024 * 1024)
agent := rpcserver.NewServer(rpcserver.ServerOptions{
ApprovedDispatcherID: approvedTestDispatcherID,
Now: func() time.Time { return clock },
Status: &agentpb.AgentStatus{
AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1", ProtocolVersion: "agent.v1",
},
})
server := grpc.NewServer()
agentpb.RegisterAgentControlServiceServer(server, agent)
go func() { _ = server.Serve(listener) }()
defer server.Stop()
conn, err := grpc.NewClient("passthrough:///bufnet", grpc.WithContextDialer(func(context.Context, string) (net.Conn, error) { return listener.Dial() }), grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
t.Fatal(err)
}
defer conn.Close()
coordinator := NewAgentCoordinator(func() time.Time { return clock })
client := agentpb.NewAgentControlServiceClient(conn)
if err := coordinator.Register("agent-1", client); err != nil {
t.Fatal(err)
}
status, err := coordinator.Probe(context.Background(), "agent-1", "cell-1")
if err != nil {
t.Fatal(err)
}
initial, err := coordinator.Activate(context.Background(), approvedTestDispatcherID, "agent-1", "cell-1", status.BootId, "epoch-1", 0)
if err != nil {
t.Fatal(err)
}
oldMeta, err := coordinator.ApprovedMeta(context.Background(), "agent-1")
if err != nil {
t.Fatal(err)
}
clock = clock.Add(6 * time.Minute)
renewed, err := coordinator.Renew(context.Background(), approvedTestDispatcherID, initial)
if err != nil {
t.Fatal(err)
}
if renewed.SessionGeneration != initial.SessionGeneration+1 || renewed.ExpiresAtUnixMs <= initial.ExpiresAtUnixMs || renewed.BootID != initial.BootID || renewed.DispatcherEpoch != initial.DispatcherEpoch {
t.Fatalf("renewal did not advance only the existing Agent session: first=%+v renewed=%+v", initial, renewed)
}
if err := coordinator.AuthorizeInboundMeta(oldMeta); err == nil {
t.Fatal("fenced pre-renewal metadata authorized an Agent report")
}
if _, err := agent.ActiveSessionMeta(); err != nil {
t.Fatalf("Agent lost its renewed session: %v", err)
}
clock = clock.Add(5 * time.Minute) // Beyond the original ten-minute expiry.
if _, err := coordinator.ApprovedMeta(context.Background(), "agent-1"); err != nil {
t.Fatalf("Agent admission expired despite timely renewal: %v", err)
}
if _, err := agent.ActiveSessionMeta(); err != nil {
t.Fatalf("Agent reports expired despite timely renewal: %v", err)
}
changed := &changedBootClient{AgentControlServiceClient: client}
if err := coordinator.Register("agent-1", changed); err != nil {
t.Fatal(err)
}
if _, err := coordinator.Renew(context.Background(), approvedTestDispatcherID, renewed); err == nil {
t.Fatal("a new Agent boot was silently adopted during renewal")
}
if changed.activated {
t.Fatal("a changed Agent boot received an activation")
}
advanced := &clockAdvancingClient{AgentControlServiceClient: client, advance: func() { clock = clock.Add(6 * time.Minute) }}
if err := coordinator.Register("agent-1", advanced); err != nil {
t.Fatal(err)
}
if _, err := coordinator.Renew(context.Background(), approvedTestDispatcherID, renewed); err == nil {
t.Fatal("a session that expired during status probing was renewed")
}
if advanced.activated {
t.Fatal("an expired session issued a fresh activation after status probing")
}
if _, err := coordinator.Renew(context.Background(), approvedTestDispatcherID, renewed); err == nil {
t.Fatal("an already expired session was renewed without explicit restart")
}
}
type changedBootClient struct {
agentpb.AgentControlServiceClient
activated bool
}
func (client *changedBootClient) GetAgentStatus(context.Context, *agentpb.GetAgentStatusRequest, ...grpc.CallOption) (*agentpb.GetAgentStatusResponse, error) {
return &agentpb.GetAgentStatusResponse{Status: &agentpb.AgentStatus{
AgentId: "agent-1", CellId: "cell-1", BootId: "different-boot", ProtocolVersion: "agent.v1",
}}, nil
}
func (client *changedBootClient) ActivateAgent(ctx context.Context, request *agentpb.ActivateAgentRequest, options ...grpc.CallOption) (*agentpb.ActivateAgentResponse, error) {
client.activated = true
return client.AgentControlServiceClient.ActivateAgent(ctx, request, options...)
}
type clockAdvancingClient struct {
agentpb.AgentControlServiceClient
advance func()
activated bool
}
func (client *clockAdvancingClient) GetAgentStatus(ctx context.Context, request *agentpb.GetAgentStatusRequest, options ...grpc.CallOption) (*agentpb.GetAgentStatusResponse, error) {
client.advance()
return client.AgentControlServiceClient.GetAgentStatus(ctx, request, options...)
}
func (client *clockAdvancingClient) ActivateAgent(ctx context.Context, request *agentpb.ActivateAgentRequest, options ...grpc.CallOption) (*agentpb.ActivateAgentResponse, error) {
client.activated = true
return client.AgentControlServiceClient.ActivateAgent(ctx, request, options...)
}