Apply approved SIP snapshots through isolated Dispatcher and Agent channel

This commit is contained in:
2026-10-03 12:45:11 +08:00
parent b68711be9c
commit a52fe5a741
32 changed files with 1437 additions and 128 deletions
+28 -19
View File
@@ -2,6 +2,7 @@ package main
import (
"bytes"
"crypto/tls"
"encoding/json"
"errors"
"fmt"
@@ -28,14 +29,6 @@ func newAgentCommand() *cobra.Command {
if err != nil {
return err
}
scenario, err := loadApprovedMockScenario(settings.MockScenarioFile)
if err != nil {
return err
}
applied, err := loadMockAppliedSIP(settings.MockAppliedSIPFile)
if err != nil {
return err
}
ca, err := readAgentPEM("MTLS_CA_FILE", settings.CAFile)
if err != nil {
return err
@@ -52,22 +45,38 @@ func newAgentCommand() *cobra.Command {
if err != nil {
return errors.New("Agent mTLS listener certificate is invalid")
}
clientTLS, err := rpc.NewClientTLSConfig(ca, cert, key, settings.DispatcherServerName)
if err != nil {
return errors.New("Agent mTLS Dispatcher certificate configuration is invalid")
var handler *rpc.Server
if mode == "sip-only" {
handler, err = newSIPOnlyAgentServer(settings)
} else {
var scenario approvedMockScenario
scenario, err = loadApprovedMockScenario(settings.MockScenarioFile)
if err != nil {
return err
}
var applied map[string]int64
applied, err = loadMockAppliedSIP(settings.MockAppliedSIPFile)
if err != nil {
return err
}
var clientTLS *tls.Config
clientTLS, err = rpc.NewClientTLSConfig(ca, cert, key, settings.DispatcherServerName)
if err != nil {
return errors.New("Agent mTLS Dispatcher certificate configuration is invalid")
}
connection, connectErr := grpc.NewClient(settings.DispatcherEndpoint, grpc.WithTransportCredentials(credentials.NewTLS(clientTLS)))
if connectErr != nil {
return errors.New("Agent cannot create pinned Dispatcher connection")
}
defer connection.Close()
handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection))
}
connection, err := grpc.NewClient(settings.DispatcherEndpoint, grpc.WithTransportCredentials(credentials.NewTLS(clientTLS)))
if err != nil {
return errors.New("Agent cannot create pinned Dispatcher connection")
}
defer connection.Close()
handler, err := newAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection))
if err != nil {
return err
}
listener, err := net.Listen("tcp", settings.Listen)
if err != nil {
return fmt.Errorf("Agent cannot listen on configured Mock address: %w", err)
return fmt.Errorf("Agent cannot listen on configured local address: %w", err)
}
defer listener.Close()
server := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS)))
@@ -91,7 +100,7 @@ func newAgentCommand() *cobra.Command {
return nil
},
}
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock mode only")
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode")
return command
}
+40
View File
@@ -0,0 +1,40 @@
package main
import (
"errors"
"os"
"path/filepath"
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
"git.ipao.vip/rogee/go-sip/internal/asterisk"
"git.ipao.vip/rogee/go-sip/internal/config"
"git.ipao.vip/rogee/go-sip/internal/rpc"
"github.com/google/uuid"
)
func newSIPOnlyAgentServer(settings config.AgentEnvironment) (*rpc.Server, error) {
if settings.AgentID == "" || settings.CellID == "" || settings.DispatcherID == "" || settings.SessionPath == "" || len(settings.PeerFingerprints) == 0 {
return nil, errors.New("SIP-only Agent requires bound identities, persistent sessions and pinned peer certificates")
}
for _, path := range []string{settings.AsteriskConfigDir, settings.AsteriskBin, settings.AsteriskLibraryDir} {
if !filepath.IsAbs(path) {
return nil, errors.New("native Asterisk paths must be explicit absolute paths")
}
}
if err := rejectLegacyAgentSpool(settings.RecoveryRoot); err != nil {
return nil, err
}
if _, err := os.Lstat(filepath.Join(settings.RecoveryRoot, ".executions")); err == nil {
return nil, errors.New("legacy Agent execution state requires operator disposition")
} else if !errors.Is(err, os.ErrNotExist) {
return nil, err
}
loader := asterisk.Loader{ConfigDir: settings.AsteriskConfigDir, Asterisk: settings.AsteriskBin, LibraryDir: settings.AsteriskLibraryDir}
return rpc.NewServer(rpc.ServerOptions{
Mode: "sip-only", StatePath: settings.SessionPath,
ApprovedDispatcherID: settings.DispatcherID,
Status: &agentpb.AgentStatus{AgentId: settings.AgentID, CellId: settings.CellID, BootId: uuid.NewString(), ProtocolVersion: "agent.v1"},
RequirePeerCertificate: true, PeerCertificateFingerprints: settings.PeerFingerprints,
LoadedSIP: loader.LoadedSIP, ApplySIP: loader.Apply,
}), nil
}
+4 -1
View File
@@ -30,7 +30,7 @@ func newDispatcherCommand() *cobra.Command {
return runDispatcher(cmd.Context(), mode)
},
}
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock mode only")
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode")
return command
}
@@ -41,6 +41,9 @@ func runDispatcher(ctx context.Context, mode string) (result error) {
if err != nil {
return err
}
if mode == "sip-only" {
return runSIPOnlyDispatcher(ctx, settings)
}
agentEndpoint, err := config.LoadMockAgentEndpoint(settings.AgentEndpointsFile)
if err != nil {
return err
+162
View File
@@ -0,0 +1,162 @@
package main
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
"git.ipao.vip/rogee/go-sip/internal/config"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
"git.ipao.vip/rogee/go-sip/internal/mq"
"git.ipao.vip/rogee/go-sip/internal/rpc"
"git.ipao.vip/rogee/go-sip/internal/store"
"github.com/google/uuid"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
)
// runSIPOnlyDispatcher owns no business queue, task discovery, recording
// listener, OSS signer, call executor or call admission. It consumes only the
// assigned control queue; unrelated controls are requeued and stop this lane.
func runSIPOnlyDispatcher(ctx context.Context, settings config.DispatcherRuntimeEnvironment) (result error) {
endpoints, err := config.LoadAgentEndpoints(settings.AgentEndpointsFile)
if err != nil {
return err
}
if len(endpoints) != 1 {
return errors.New("SIP-only Dispatcher requires exactly one assigned Agent")
}
endpoint := endpoints[0]
ca, err := readAgentPEM("MTLS_CA_FILE", settings.CAFile)
if err != nil {
return err
}
cert, err := readAgentPEM("MTLS_CERT_FILE", settings.CertFile)
if err != nil {
return err
}
key, err := readAgentPEM("MTLS_KEY_FILE", settings.KeyFile)
if err != nil {
return err
}
clientTLS, err := rpc.NewClientTLSConfig(ca, cert, key, endpoint.ServerName)
if err != nil {
return errors.New("SIP-only Dispatcher mTLS Agent certificate is invalid")
}
httpClient, err := localMockHTTPClient(ca)
if err != nil {
return err
}
httpClient.Timeout = 10 * time.Second
_, reader, err := dispatcherConfigurationClient("sip-only", httpClient)
if err != nil {
return err
}
db, err := store.Open(settings.SQLitePath)
if err != nil {
return err
}
defer func() { result = errors.Join(result, db.Close()) }()
if err := db.CloseAdmission(settings.DispatcherID); err != nil {
return err
}
broker, err := mq.Open(settings.RabbitMQURL, settings.DispatcherID, 1)
if err != nil {
return err
}
defer func() { result = errors.Join(result, broker.Close()) }()
connection, err := grpc.NewClient(endpoint.Address, grpc.WithTransportCredentials(credentials.NewTLS(clientTLS)))
if err != nil {
return errors.New("SIP-only Dispatcher cannot create pinned Agent connection")
}
defer func() { result = errors.Join(result, connection.Close()) }()
agent := agentpb.NewAgentControlServiceClient(connection)
coordinator := dispatcher.NewAgentCoordinator(time.Now)
if err := coordinator.Register(endpoint.AgentID, agent); err != nil {
return err
}
probeCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
status, err := coordinator.Probe(probeCtx, endpoint.AgentID, endpoint.CellID)
if err != nil {
cancel()
return fmt.Errorf("probe SIP-only assigned Agent: %w", err)
}
epoch, err := uuid.NewRandom()
if err != nil {
cancel()
return err
}
session, err := coordinator.Activate(probeCtx, settings.DispatcherID, endpoint.AgentID, endpoint.CellID, status.BootId, epoch.String(), 0)
cancel()
if err != nil {
return fmt.Errorf("activate SIP-only Agent session: %w", err)
}
originator := &dispatcher.ApprovedOriginator{DispatcherID: settings.DispatcherID, Client: agent, Meta: func(callCtx context.Context) (*agentpb.RequestMeta, error) {
return coordinator.ApprovedMeta(callCtx, endpoint.AgentID)
}}
lane := dispatcher.SIPOnly{DispatcherID: settings.DispatcherID, Store: db, Client: reader, ApplySIP: originator.ApplySIP, VerifySIP: originator.VerifySIP}
if _, err := broker.DrainControlPredeclared(ctx, broker.ControlQueue(), lane.HandleNotification); err != nil {
return fmt.Errorf("drain assigned SIP controls before startup: %w", err)
}
if err := lane.Sync(ctx); err != nil {
return fmt.Errorf("initial native SIP-only sync: %w", err)
}
serveCtx, stop := context.WithCancel(ctx)
defer stop()
renewErrors := make(chan error, 1)
go func() {
renewErr := maintainAgentSession(serveCtx, session, 5*time.Minute, func(c context.Context, previous dispatcher.AgentSession) (dispatcher.AgentSession, error) {
return coordinator.Renew(c, settings.DispatcherID, previous)
})
renewErrors <- renewErr
stop()
}()
consumer, err := broker.StartPredeclaredConsumer(serveCtx, broker.ControlQueue(), lane.HandleNotification)
if err != nil {
return err
}
defer func() {
waitCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
result = errors.Join(result, consumer.Stop(waitCtx))
}()
consumerErrors := make(chan error, 1)
go func() { consumerErrors <- consumer.Wait(context.Background()); stop() }()
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return nil
case err := <-renewErrors:
if err != nil {
return fmt.Errorf("SIP-only Agent session renewal failed: %w", err)
}
return errors.New("SIP-only Agent session stopped unexpectedly")
case err := <-consumerErrors:
if err != nil {
return fmt.Errorf("SIP-only control queue stopped: %w", err)
}
return errors.New("SIP-only control queue stopped unexpectedly")
case err := <-broker.Done():
return fmt.Errorf("SIP-only broker connection lost: %v", err)
case <-ticker.C:
_, pending, err := db.SIPState(settings.DispatcherID)
if err != nil {
return err
}
if pending == 0 {
continue
}
readCtx, cancel := context.WithTimeout(serveCtx, 20*time.Second)
err = lane.Sync(readCtx)
cancel()
if err != nil {
slog.Warn("SIP-only revision remains durably pending", "dispatcher_id", settings.DispatcherID, "pending_revision", pending, "error", err)
}
}
}
}