feat(nonprod): bind native Agent to signed calls and live capture
This commit is contained in:
@@ -22,7 +22,7 @@ import (
|
||||
func newAgentCommand() *cobra.Command {
|
||||
var mode string
|
||||
command := &cobra.Command{
|
||||
Use: "agent", Short: "Run the isolated Agent", Args: cobra.NoArgs,
|
||||
Use: "agent", Short: "Run the bound Agent", Args: cobra.NoArgs,
|
||||
SilenceUsage: true,
|
||||
RunE: func(cmd *cobra.Command, _ []string) error {
|
||||
settings, err := config.LoadAgentEnvironment(mode)
|
||||
@@ -50,14 +50,16 @@ func newAgentCommand() *cobra.Command {
|
||||
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
|
||||
if mode == "mock" {
|
||||
scenario, err = loadApprovedMockScenario(settings.MockScenarioFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
applied, err = loadMockAppliedSIP(settings.MockAppliedSIPFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
var clientTLS *tls.Config
|
||||
clientTLS, err = rpc.NewClientTLSConfig(ca, cert, key, settings.DispatcherServerName)
|
||||
@@ -69,7 +71,12 @@ func newAgentCommand() *cobra.Command {
|
||||
return errors.New("Agent cannot create pinned Dispatcher connection")
|
||||
}
|
||||
defer connection.Close()
|
||||
handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection))
|
||||
client := agentpb.NewAgentControlServiceClient(connection)
|
||||
if mode == "mock" {
|
||||
handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, client)
|
||||
} else {
|
||||
handler, err = newRealAgentServer(cmd.Context(), settings, client)
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -100,12 +107,12 @@ func newAgentCommand() *cobra.Command {
|
||||
return nil
|
||||
},
|
||||
}
|
||||
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode")
|
||||
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock, SIP-only or explicit nonprod-real mode")
|
||||
return command
|
||||
}
|
||||
|
||||
// This is an isolated deployment fixture, not SaaS SIP configuration or proof
|
||||
// that Asterisk actually loaded the revisions. Real/mixed startup is rejected.
|
||||
// that Asterisk actually loaded the revisions. Nonprod-real cannot load it.
|
||||
func loadMockAppliedSIP(path string) (map[string]int64, error) {
|
||||
if strings.TrimSpace(path) == "" {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE is required")
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"reflect"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/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"
|
||||
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
// newRealAgentServer only exposes the signed Dispatcher execution path. There
|
||||
// is no scenario, synthetic recording, local dial entry, or Mock upload host.
|
||||
func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, dispatcher agentpb.AgentControlServiceClient) (*rpc.Server, error) {
|
||||
if ctx == nil || ctx.Err() != nil || settings.Mode != "nonprod-real" || settings.AgentID == "" || settings.CellID == "" ||
|
||||
settings.SessionPath == "" || settings.RecoveryRoot == "" || len(settings.PeerFingerprints) == 0 ||
|
||||
settings.OSSAllowedHost == "" || dispatcher == nil || tenant.ValidateDispatcherID(settings.DispatcherID) != nil {
|
||||
return nil, errors.New("real Agent requires approved nonproduction identity, pinned Dispatcher and private recovery")
|
||||
}
|
||||
value := reflect.ValueOf(dispatcher)
|
||||
if value.Kind() == reflect.Pointer && value.IsNil() {
|
||||
return nil, errors.New("real Agent requires a live pinned Dispatcher transport")
|
||||
}
|
||||
for _, path := range []string{settings.AsteriskConfigDir, settings.AsteriskBin, settings.AsteriskLibraryDir, settings.EvidenceRoot, settings.RecoveryRoot} {
|
||||
if !filepath.IsAbs(path) {
|
||||
return nil, errors.New("real Agent paths must be absolute")
|
||||
}
|
||||
}
|
||||
if strings.ContainsAny(settings.OSSAllowedHost, "/@ ") || settings.MockScenarioFile != "" || settings.MockAppliedSIPFile != "" {
|
||||
return nil, errors.New("real Agent must not use Mock assets or an invalid OSS host")
|
||||
}
|
||||
root, err := os.Stat(settings.RecoveryRoot)
|
||||
if err != nil || !root.IsDir() || root.Mode().Perm() != 0700 {
|
||||
return nil, errors.New("real Agent recovery directory must exist with mode 0700")
|
||||
}
|
||||
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}
|
||||
var handler *rpc.Server
|
||||
worker := &rpc.ApprovedCallWorker{
|
||||
Lifecycle: ctx, Calls: &agent.TaskCalls{},
|
||||
Prepare: func(execution rpc.ApprovedExecution) (func(context.Context) error, error) {
|
||||
if handler == nil {
|
||||
return nil, errors.New("Agent session is unavailable")
|
||||
}
|
||||
delivery := &agent.RecordingDelivery{
|
||||
Call: agent.RecordingClient{
|
||||
Client: dispatcher, DispatcherID: execution.DispatcherID, TenantID: execution.TenantID,
|
||||
SourceEventID: execution.SourceEventID, Session: func(context.Context) (*agentpb.RequestMeta, error) { return handler.ActiveSessionMeta() },
|
||||
},
|
||||
Recovery: &agent.RecordingRecovery{
|
||||
Root: settings.RecoveryRoot,
|
||||
Upload: agent.UploadClient{AllowedHosts: map[string]struct{}{strings.ToLower(settings.OSSAllowedHost): {}}},
|
||||
},
|
||||
}
|
||||
return (&rpc.ApprovedRecordedRealCall{
|
||||
Loader: loader, MediaPayloadType: 118, EvidenceRoot: settings.EvidenceRoot,
|
||||
MaxWAVBytes: 64 << 20, ReportTimeout: 15 * time.Minute, Delivery: delivery,
|
||||
}).Prepare(execution)
|
||||
},
|
||||
OnFailure: func(execution rpc.ApprovedExecution, cause error) error {
|
||||
log.Printf("Agent real call requires inspection: event_id=%q task_id=%q cause_type=%T", execution.SourceEventID, execution.TaskID, cause)
|
||||
return nil
|
||||
},
|
||||
}
|
||||
handler, err = rpc.NewApprovedAgentServer(rpc.ServerOptions{
|
||||
Mode: settings.Mode, 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,
|
||||
}, worker)
|
||||
return handler, err
|
||||
}
|
||||
@@ -0,0 +1,116 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/rpc"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials"
|
||||
)
|
||||
|
||||
func TestRealAgentCommandStartsPinnedServerWithoutMockFixturesOrDial(t *testing.T) {
|
||||
ca, agentCert, agentKey, dispatcherCert, dispatcherKey, dispatcherLeaf := localCommandCertificates(t)
|
||||
settings, _ := currentAgentSetupFixture(t)
|
||||
for path, data := range map[string][]byte{settings.CAFile: ca, settings.CertFile: agentCert, settings.KeyFile: agentKey} {
|
||||
if err := os.WriteFile(path, data, 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
listen := listener.Addr().String()
|
||||
listener.Close()
|
||||
for name, value := range map[string]string{
|
||||
"AGENT_ID": settings.AgentID, "CELL_ID": settings.CellID,
|
||||
"AGENT_GRPC_LISTEN": listen, "AGENT_SESSION_PATH": settings.SessionPath,
|
||||
"AGENT_RECOVERY_ROOT": settings.RecoveryRoot, "DISPATCHER_ID": settings.DispatcherID,
|
||||
"DISPATCHER_GRPC_ENDPOINT": settings.DispatcherEndpoint, "DISPATCHER_GRPC_SERVER_NAME": settings.DispatcherServerName,
|
||||
"MTLS_CA_FILE": settings.CAFile, "MTLS_CERT_FILE": settings.CertFile, "MTLS_KEY_FILE": settings.KeyFile,
|
||||
"MTLS_PEER_CERT_FINGERPRINTS": rpc.CertificateFingerprint(dispatcherLeaf),
|
||||
"ASTERISK_CONFIG_DIR": filepath.Join(settings.RecoveryRoot, "asterisk"),
|
||||
"ASTERISK_BIN": "/usr/sbin/asterisk", "ASTERISK_LIBRARY_DIR": "/usr/lib/asterisk/modules",
|
||||
"AGENT_EVIDENCE_ROOT": filepath.Join(settings.RecoveryRoot, "evidence"),
|
||||
"AGENT_OSS_ALLOWED_HOST": "bucket.oss-cn-beijing.aliyuncs.com",
|
||||
"AGENT_MOCK_SCENARIO_FILE": "", "AGENT_MOCK_APPLIED_SIP_FILE": "",
|
||||
} {
|
||||
t.Setenv(name, value)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cmd := newAgentCommand()
|
||||
cmd.SetOut(io.Discard)
|
||||
cmd.SetErr(io.Discard)
|
||||
cmd.SetContext(ctx)
|
||||
cmd.SetArgs([]string{"--mode", "nonprod-real"})
|
||||
finished := make(chan error, 1)
|
||||
go func() { finished <- cmd.Execute() }()
|
||||
t.Cleanup(func() {
|
||||
cancel()
|
||||
select {
|
||||
case err := <-finished:
|
||||
if err != nil {
|
||||
t.Errorf("real Agent shutdown failed: %v", err)
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Error("real Agent did not stop after cancellation")
|
||||
}
|
||||
})
|
||||
tlsConfig, err := rpc.NewClientTLSConfig(ca, dispatcherCert, dispatcherKey, "agent.local")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
dialCtx, stop := context.WithTimeout(context.Background(), 6*time.Second)
|
||||
defer stop()
|
||||
conn, err := grpc.DialContext(dialCtx, listen, grpc.WithTransportCredentials(credentials.NewTLS(tlsConfig)), grpc.WithBlock())
|
||||
if err != nil {
|
||||
select {
|
||||
case startupErr := <-finished:
|
||||
t.Fatalf("nonprod-real startup failed: %v (dial: %v)", startupErr, err)
|
||||
default:
|
||||
t.Fatalf("nonprod-real pinned server unavailable: %v", err)
|
||||
}
|
||||
}
|
||||
defer conn.Close()
|
||||
status, err := agentpb.NewAgentControlServiceClient(conn).GetAgentStatus(dialCtx, &agentpb.GetAgentStatusRequest{
|
||||
Meta: &agentpb.RequestMeta{ProtocolVersion: "agent.v1", RequestId: "status", TraceId: "status", OperationId: "status", AgentId: settings.AgentID, CellId: settings.CellID},
|
||||
Target: &agentpb.AgentBinding{AgentId: settings.AgentID, CellId: settings.CellID},
|
||||
})
|
||||
if err != nil || status.GetStatus().GetAgentId() != settings.AgentID {
|
||||
t.Fatalf("nonprod-real Agent did not expose its pinned status: status=%v err=%v", status, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRealAgentServerRequiresBoundNativeSIPAndNoMock(t *testing.T) {
|
||||
settings, _ := currentAgentSetupFixture(t)
|
||||
settings.Mode = "nonprod-real"
|
||||
settings.AsteriskConfigDir = filepath.Join(settings.RecoveryRoot, "asterisk")
|
||||
settings.AsteriskBin = "/usr/sbin/asterisk"
|
||||
settings.AsteriskLibraryDir = "/usr/lib/asterisk/modules"
|
||||
settings.EvidenceRoot = filepath.Join(settings.RecoveryRoot, "evidence")
|
||||
settings.OSSAllowedHost = "bucket.oss-cn-beijing.aliyuncs.com"
|
||||
client := &isolatedAgentRecordingClient{}
|
||||
server, err := newRealAgentServer(context.Background(), settings, client)
|
||||
if err != nil || server == nil {
|
||||
t.Fatalf("real Agent must use the native SIP and non-Mock delivery boundary: %v", err)
|
||||
}
|
||||
if _, err := os.Lstat(settings.SessionPath); !os.IsNotExist(err) {
|
||||
t.Fatalf("server assembly must not touch the existing session: %v", err)
|
||||
}
|
||||
var typedNil *isolatedAgentRecordingClient
|
||||
var nilClient agentpb.AgentControlServiceClient = typedNil
|
||||
if server, err := newRealAgentServer(context.Background(), settings, nilClient); err == nil || server != nil {
|
||||
t.Fatal("typed-nil Dispatcher must block real Agent startup")
|
||||
}
|
||||
settings.MockScenarioFile = "old-mock.json"
|
||||
if server, err := newRealAgentServer(context.Background(), settings, client); err == nil || server != nil {
|
||||
t.Fatal("real Agent must not accept isolated Mock fixtures")
|
||||
}
|
||||
}
|
||||
@@ -24,13 +24,13 @@ import (
|
||||
func newDispatcherCommand() *cobra.Command {
|
||||
var mode string
|
||||
command := &cobra.Command{
|
||||
Use: "dispatcher", Short: "Run the isolated Dispatcher", Args: cobra.NoArgs,
|
||||
Use: "dispatcher", Short: "Run the bound Dispatcher", Args: cobra.NoArgs,
|
||||
SilenceUsage: true,
|
||||
RunE: func(cmd *cobra.Command, _ []string) error {
|
||||
return runDispatcher(cmd.Context(), mode)
|
||||
},
|
||||
}
|
||||
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode")
|
||||
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock, SIP-only or explicit nonprod-real mode")
|
||||
return command
|
||||
}
|
||||
|
||||
@@ -48,7 +48,11 @@ func runDispatcher(ctx context.Context, mode string) (result error) {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
ossConfiguration, err := config.LoadOSSConfig(settings.OSSConfigFile, settings.DispatcherID)
|
||||
loadOSS := config.LoadOSSConfig
|
||||
if mode == "nonprod-real" {
|
||||
loadOSS = config.LoadNonprodRealOSSConfig
|
||||
}
|
||||
ossConfiguration, err := loadOSS(settings.OSSConfigFile, settings.DispatcherID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -147,7 +151,7 @@ func runDispatcher(ctx context.Context, mode string) (result error) {
|
||||
}
|
||||
listener, err := net.Listen("tcp", settings.Listen)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Dispatcher cannot listen on configured Mock address: %w", err)
|
||||
return fmt.Errorf("Dispatcher cannot listen on configured local address: %w", err)
|
||||
}
|
||||
server := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS)))
|
||||
agentpb.RegisterAgentControlServiceServer(server, &rpc.RecordingServer{
|
||||
|
||||
Reference in New Issue
Block a user