227 lines
9.3 KiB
Go
227 lines
9.3 KiB
Go
package dispatcher
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"sync"
|
|
"time"
|
|
|
|
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
|
agentsession "git.ipao.vip/rogee/go-sip/internal/session"
|
|
"git.ipao.vip/rogee/go-sip/internal/store"
|
|
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
|
"github.com/google/uuid"
|
|
)
|
|
|
|
type AgentSession struct {
|
|
AgentID string
|
|
CellID string
|
|
BootID string
|
|
DispatcherEpoch string
|
|
SessionGeneration uint64
|
|
ExpiresAtUnixMs int64
|
|
}
|
|
|
|
// AgentCoordinator binds an Agent endpoint to the current activated session.
|
|
type AgentCoordinator struct {
|
|
mu sync.Mutex
|
|
now func() time.Time
|
|
clients map[string]agentpb.AgentControlServiceClient
|
|
sessions map[string]AgentSession
|
|
}
|
|
|
|
func NewAgentCoordinator(now func() time.Time) *AgentCoordinator {
|
|
if now == nil {
|
|
now = time.Now
|
|
}
|
|
return &AgentCoordinator{now: now, clients: make(map[string]agentpb.AgentControlServiceClient), sessions: make(map[string]AgentSession)}
|
|
}
|
|
|
|
func (c *AgentCoordinator) Register(agentID string, client agentpb.AgentControlServiceClient) error {
|
|
if agentID == "" || client == nil {
|
|
return errors.New("agent ID and client are required")
|
|
}
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
c.clients[agentID] = client
|
|
return nil
|
|
}
|
|
|
|
// Probe reads deployment status before Activate binds the returned boot ID.
|
|
func (c *AgentCoordinator) Probe(ctx context.Context, agentID, cellID string) (*agentpb.AgentStatus, error) {
|
|
if agentID == "" || cellID == "" {
|
|
return nil, errors.New("agent ID and cell ID are required")
|
|
}
|
|
client, err := c.client(agentID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
operationID := "status:" + agentID
|
|
response, err := client.GetAgentStatus(ctx, &agentpb.GetAgentStatusRequest{
|
|
Meta: &agentpb.RequestMeta{
|
|
ProtocolVersion: "agent.v1",
|
|
RequestId: operationID + ":request",
|
|
TraceId: operationID,
|
|
OperationId: operationID,
|
|
AgentId: agentID,
|
|
CellId: cellID,
|
|
},
|
|
Target: &agentpb.AgentBinding{AgentId: agentID, CellId: cellID},
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if response == nil || response.Status == nil {
|
|
return nil, errors.New("Agent status response is empty")
|
|
}
|
|
if response.Status.AgentId != agentID || response.Status.CellId != cellID || response.Status.BootId == "" {
|
|
return nil, errors.New("Agent status identity or boot ID is invalid")
|
|
}
|
|
return response.Status, nil
|
|
}
|
|
|
|
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
|
|
}
|
|
if cellID == "" || bootID == "" || epoch == "" {
|
|
return AgentSession{}, errors.New("cell, boot and dispatcher epoch are required")
|
|
}
|
|
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, DispatcherId: dispatcherID, DispatcherEpoch: epoch, SessionGeneration: generation}, ActivationOperationId: operationID})
|
|
if err != nil {
|
|
return AgentSession{}, err
|
|
}
|
|
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()
|
|
c.sessions[agentID] = session
|
|
c.mu.Unlock()
|
|
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()
|
|
client := c.clients[agentID]
|
|
if client == nil {
|
|
return nil, fmt.Errorf("Agent %q is not registered", agentID)
|
|
}
|
|
return client, nil
|
|
}
|
|
|
|
func (c *AgentCoordinator) clientAndSession(agentID string) (agentpb.AgentControlServiceClient, AgentSession, error) {
|
|
client, err := c.client(agentID)
|
|
if err != nil {
|
|
return nil, AgentSession{}, err
|
|
}
|
|
c.mu.Lock()
|
|
session, ok := c.sessions[agentID]
|
|
c.mu.Unlock()
|
|
if !ok {
|
|
return nil, AgentSession{}, fmt.Errorf("Agent %q is not activated", agentID)
|
|
}
|
|
return client, session, nil
|
|
}
|
|
|
|
// ApprovedMeta supplies a fresh operation identity from the current session.
|
|
func (c *AgentCoordinator) ApprovedMeta(ctx context.Context, agentID string) (*agentpb.RequestMeta, error) {
|
|
if c == nil || ctx == nil || agentID == "" {
|
|
return nil, errors.New("registered Agent and request context are required")
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
_, session, err := c.clientAndSession(agentID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if session.AgentID != agentID || session.CellID == "" || session.BootID == "" || session.DispatcherEpoch == "" ||
|
|
session.SessionGeneration == 0 || session.ExpiresAtUnixMs <= c.now().UnixMilli() {
|
|
return nil, errors.New("approved Agent session is incomplete or expired")
|
|
}
|
|
identity, err := uuid.NewRandom()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create Agent operation identity: %w", err)
|
|
}
|
|
return c.meta(session, "approved:"+identity.String()), nil
|
|
}
|
|
|
|
// AuthorizeInboundMeta accepts an Agent fact only from the active Endpoint
|
|
// session; a replayed fact must carry the newly activated boot and generation.
|
|
func (c *AgentCoordinator) AuthorizeInboundMeta(meta *agentpb.RequestMeta) error {
|
|
if c == nil || meta == nil || meta.ProtocolVersion != "agent.v1" || meta.AgentId == "" {
|
|
return fmt.Errorf("active Agent session is required: %w", store.ErrCommandConflict)
|
|
}
|
|
c.mu.Lock()
|
|
session, active := c.sessions[meta.AgentId]
|
|
client, registered := c.clients[meta.AgentId]
|
|
c.mu.Unlock()
|
|
if !active || !registered || client == nil || session.ExpiresAtUnixMs <= c.now().UnixMilli() ||
|
|
meta.CellId != session.CellID || meta.BootId != session.BootID ||
|
|
meta.DispatcherEpoch != session.DispatcherEpoch {
|
|
return fmt.Errorf("Agent report does not match a current activated session: %w", store.ErrCommandConflict)
|
|
}
|
|
if meta.SessionGeneration != session.SessionGeneration {
|
|
if meta.SessionGeneration > 0 &&
|
|
((meta.SessionGeneration < session.SessionGeneration && session.SessionGeneration-meta.SessionGeneration == 1) ||
|
|
(meta.SessionGeneration > session.SessionGeneration && meta.SessionGeneration-session.SessionGeneration == 1)) {
|
|
return fmt.Errorf("Agent fact was fenced during session transition: %w: %w", store.ErrCommandConflict, agentsession.ErrGenerationTransition)
|
|
}
|
|
return fmt.Errorf("Agent report has an unapproved session generation: %w", store.ErrCommandConflict)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (c *AgentCoordinator) meta(session AgentSession, operationID string) *agentpb.RequestMeta {
|
|
return &agentpb.RequestMeta{ProtocolVersion: "agent.v1", RequestId: operationID + ":request", TraceId: operationID, OperationId: operationID, DispatcherEpoch: session.DispatcherEpoch, AgentId: session.AgentID, CellId: session.CellID, BootId: session.BootID, SessionGeneration: session.SessionGeneration}
|
|
}
|