Route Agent CLI through isolated approved Mock runtime
This commit is contained in:
@@ -0,0 +1,146 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/config"
|
||||
"git.ipao.vip/rogee/go-sip/internal/rpc"
|
||||
"github.com/spf13/cobra"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials"
|
||||
)
|
||||
|
||||
func newCurrentAgentCommand() *cobra.Command {
|
||||
var mode string
|
||||
command := &cobra.Command{
|
||||
Use: "agent", Short: "Run the isolated Agent", Args: cobra.NoArgs,
|
||||
SilenceUsage: true,
|
||||
RunE: func(cmd *cobra.Command, _ []string) error {
|
||||
settings, err := config.LoadAgentEnvironment(mode)
|
||||
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
|
||||
}
|
||||
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
|
||||
}
|
||||
serverTLS, err := rpc.NewServerTLSConfig(ca, cert, key)
|
||||
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")
|
||||
}
|
||||
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 := newCurrentAgentServer(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)
|
||||
}
|
||||
defer listener.Close()
|
||||
server := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS)))
|
||||
agentpb.RegisterAgentControlServiceServer(server, handler)
|
||||
stopped := make(chan struct{})
|
||||
go func() {
|
||||
select {
|
||||
case <-cmd.Context().Done():
|
||||
server.GracefulStop()
|
||||
case <-stopped:
|
||||
}
|
||||
}()
|
||||
err = server.Serve(listener)
|
||||
close(stopped)
|
||||
if err != nil && !errors.Is(err, grpc.ErrServerStopped) {
|
||||
return fmt.Errorf("Agent gRPC listener stopped: %w", err)
|
||||
}
|
||||
if cmd.Context().Err() == nil {
|
||||
return errors.New("Agent gRPC listener stopped without shutdown")
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock mode only")
|
||||
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.
|
||||
func loadMockAppliedSIP(path string) (map[string]int64, error) {
|
||||
if strings.TrimSpace(path) == "" {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE is required")
|
||||
}
|
||||
file, err := os.Open(path)
|
||||
if err != nil {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE cannot be opened")
|
||||
}
|
||||
defer file.Close()
|
||||
data, err := io.ReadAll(io.LimitReader(file, 64<<10+1))
|
||||
if err != nil || len(data) > 64<<10 {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE exceeds its size limit or cannot be read")
|
||||
}
|
||||
var fixture struct {
|
||||
Trunks map[string]int64 `json:"trunks"`
|
||||
}
|
||||
decoder := json.NewDecoder(bytes.NewReader(data))
|
||||
decoder.DisallowUnknownFields()
|
||||
if err := decoder.Decode(&fixture); err != nil {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE has invalid fields")
|
||||
}
|
||||
if err := decoder.Decode(new(any)); err != io.EOF {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE contains extra data")
|
||||
}
|
||||
if len(fixture.Trunks) == 0 {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE requires explicit trunks")
|
||||
}
|
||||
for trunk, revision := range fixture.Trunks {
|
||||
if strings.TrimSpace(trunk) == "" || revision <= 0 {
|
||||
return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE has invalid trunk revision")
|
||||
}
|
||||
}
|
||||
return fixture.Trunks, nil
|
||||
}
|
||||
|
||||
func readAgentPEM(name, path string) ([]byte, error) {
|
||||
file, err := os.Open(path)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s cannot be opened", name)
|
||||
}
|
||||
defer file.Close()
|
||||
data, err := io.ReadAll(io.LimitReader(file, 1<<20+1))
|
||||
if err != nil || len(data) == 0 || len(data) > 1<<20 {
|
||||
return nil, fmt.Errorf("%s is unreadable or too large", name)
|
||||
}
|
||||
return data, nil
|
||||
}
|
||||
@@ -0,0 +1,168 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"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/codes"
|
||||
"google.golang.org/grpc/credentials"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
func TestCurrentAgentCommandRejectsRealModesAndObsoleteCallEntry(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
args []string
|
||||
want string
|
||||
}{
|
||||
{"real", []string{"--mode", "real"}, "Mock"},
|
||||
{"mixed", []string{"--mode", "mixed"}, "Mock"},
|
||||
{"old call-once", []string{"--call-once"}, "unknown flag"},
|
||||
{"no deployment identity", []string{"--mode", "mock"}, "AGENT_ID"},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
state := filepath.Join(t.TempDir(), "session.json")
|
||||
t.Setenv("AGENT_ID", "")
|
||||
t.Setenv("AGENT_SESSION_PATH", state)
|
||||
command := newCurrentAgentCommand()
|
||||
command.SetOut(io.Discard)
|
||||
command.SetErr(io.Discard)
|
||||
command.SetArgs(tc.args)
|
||||
if err := command.Execute(); err == nil || !strings.Contains(err.Error(), tc.want) {
|
||||
t.Fatalf("unsafe current Agent command was admitted: %v", err)
|
||||
}
|
||||
if _, err := os.Stat(state); !os.IsNotExist(err) {
|
||||
t.Fatalf("rejected command opened durable Agent state: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadMockAppliedSIPAcceptsOnlyExplicitLocalFacts(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "applied-sip.json")
|
||||
if err := os.WriteFile(path, []byte(`{"trunks":{"trunk-mock":8}}`), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
applied, err := loadMockAppliedSIP(path)
|
||||
if err != nil || len(applied) != 1 || applied["trunk-mock"] != 8 {
|
||||
t.Fatalf("explicit isolated applied SIP facts were changed: %v %v", applied, err)
|
||||
}
|
||||
if _, err := loadMockAppliedSIP(""); err == nil {
|
||||
t.Fatal("missing applied SIP facts were defaulted")
|
||||
}
|
||||
for _, data := range []string{
|
||||
`{"trunks":{}}`, `{"trunks":{"trunk-mock":0}}`, `{"trunks":{"":8}}`,
|
||||
`{"trunks":{"trunk-mock":8},"schema_version":"v3"}`,
|
||||
`{"trunks":{"trunk-mock":8}} {}`,
|
||||
} {
|
||||
if err := os.WriteFile(path, []byte(data), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := loadMockAppliedSIP(path); err == nil {
|
||||
t.Fatalf("unapproved Mock SIP fixture was admitted: %s", data)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentAgentCommandServesPinnedMutualTLSOnly(t *testing.T) {
|
||||
ca, agentCert, agentKey, dispatcherCert, dispatcherKey, dispatcherLeaf := localCommandCertificates(t)
|
||||
settings, _ := currentAgentSetupFixture(t)
|
||||
for path, body := range map[string][]byte{
|
||||
settings.CAFile: ca, settings.CertFile: agentCert, settings.KeyFile: agentKey,
|
||||
} {
|
||||
if err := os.WriteFile(path, body, 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
scenarioFile := filepath.Join(settings.RecoveryRoot, "scenario.json")
|
||||
appliedFile := filepath.Join(settings.RecoveryRoot, "applied-sip.json")
|
||||
if err := os.WriteFile(scenarioFile, []byte(`{"inbound_pcm16":"AQABAA==","script":{"Turns":[{"Transcript":"synthetic fixture"}]},"max_wav_bytes":4096,"expected_recording":false,"outcome":"no_answer","reason_message":"synthetic no answer"}`), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(appliedFile, []byte(`{"trunks":{"trunk-mock":8}}`), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
reserved, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
settings.Listen = reserved.Addr().String()
|
||||
_ = reserved.Close()
|
||||
for name, value := range map[string]string{
|
||||
"AGENT_ID": settings.AgentID, "CELL_ID": settings.CellID,
|
||||
"AGENT_GRPC_LISTEN": settings.Listen, "AGENT_SESSION_PATH": settings.SessionPath,
|
||||
"AGENT_RECOVERY_ROOT": settings.RecoveryRoot,
|
||||
"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),
|
||||
"AGENT_MOCK_SCENARIO_FILE": scenarioFile, "AGENT_MOCK_APPLIED_SIP_FILE": appliedFile,
|
||||
} {
|
||||
t.Setenv(name, value)
|
||||
}
|
||||
process, cancel := context.WithCancel(context.Background())
|
||||
command := newCurrentAgentCommand()
|
||||
command.SetOut(io.Discard)
|
||||
command.SetErr(io.Discard)
|
||||
command.SetContext(process)
|
||||
command.SetArgs([]string{"--mode", "mock"})
|
||||
finished := make(chan error, 1)
|
||||
go func() { finished <- command.Execute() }()
|
||||
t.Cleanup(func() {
|
||||
cancel()
|
||||
select {
|
||||
case err := <-finished:
|
||||
if err != nil {
|
||||
t.Errorf("isolated Agent shutdown failed: %v", err)
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Error("isolated Agent did not stop after process cancellation")
|
||||
}
|
||||
})
|
||||
trustedTLS, err := rpc.NewClientTLSConfig(ca, dispatcherCert, dispatcherKey, "agent.local")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
dial, stopDial := context.WithTimeout(context.Background(), 6*time.Second)
|
||||
defer stopDial()
|
||||
trusted, err := grpc.DialContext(dial, settings.Listen, grpc.WithTransportCredentials(credentials.NewTLS(trustedTLS)), grpc.WithBlock())
|
||||
if err != nil {
|
||||
select {
|
||||
case startupErr := <-finished:
|
||||
t.Fatalf("Agent startup failed: %v (dial: %v)", startupErr, err)
|
||||
default:
|
||||
t.Fatalf("isolated Agent mTLS listener unavailable: %v", err)
|
||||
}
|
||||
}
|
||||
defer trusted.Close()
|
||||
probe := &agentpb.GetAgentStatusRequest{
|
||||
Meta: &agentpb.RequestMeta{ProtocolVersion: "agent.v1", RequestId: "mock-status", TraceId: "mock-status", OperationId: "mock-status", AgentId: settings.AgentID, CellId: settings.CellID},
|
||||
Target: &agentpb.AgentBinding{AgentId: settings.AgentID, CellId: settings.CellID},
|
||||
}
|
||||
response, err := agentpb.NewAgentControlServiceClient(trusted).GetAgentStatus(dial, probe)
|
||||
if err != nil || response.GetStatus().GetAgentId() != settings.AgentID || response.GetStatus().GetCellId() != settings.CellID {
|
||||
t.Fatalf("trusted D could not probe the actual Agent command: response=%v err=%v", response, err)
|
||||
}
|
||||
untrustedTLS, err := rpc.NewClientTLSConfig(ca, agentCert, agentKey, "agent.local")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
untrusted, err := grpc.NewClient(settings.Listen, grpc.WithTransportCredentials(credentials.NewTLS(untrustedTLS)))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer untrusted.Close()
|
||||
if _, err := agentpb.NewAgentControlServiceClient(untrusted).GetAgentStatus(dial, probe); status.Code(err) != codes.PermissionDenied {
|
||||
t.Fatalf("unpinned but CA-valid client could probe Agent: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"crypto/ecdsa"
|
||||
"crypto/elliptic"
|
||||
"crypto/rand"
|
||||
"crypto/x509"
|
||||
"crypto/x509/pkix"
|
||||
"encoding/pem"
|
||||
"math/big"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// localCommandCertificates creates isolated fixtures for both command-line
|
||||
// roles. No certificate or key is stored outside the test's temporary files.
|
||||
func localCommandCertificates(t *testing.T) (ca, agentCert, agentKey, dispatcherCert, dispatcherKey []byte, dispatcherLeaf *x509.Certificate) {
|
||||
t.Helper()
|
||||
caKey, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
now := time.Now()
|
||||
root := &x509.Certificate{
|
||||
SerialNumber: big.NewInt(1), Subject: pkix.Name{CommonName: "isolated test root"},
|
||||
NotBefore: now.Add(-time.Minute), NotAfter: now.Add(time.Hour),
|
||||
IsCA: true, BasicConstraintsValid: true, KeyUsage: x509.KeyUsageCertSign | x509.KeyUsageDigitalSignature,
|
||||
}
|
||||
rootDER, err := x509.CreateCertificate(rand.Reader, root, root, &caKey.PublicKey, caKey)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
root, err = x509.ParseCertificate(rootDER)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ca = pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: rootDER})
|
||||
issue := func(serial int64, dns string) ([]byte, []byte, *x509.Certificate) {
|
||||
t.Helper()
|
||||
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
template := &x509.Certificate{
|
||||
SerialNumber: big.NewInt(serial), Subject: pkix.Name{CommonName: dns},
|
||||
NotBefore: now.Add(-time.Minute), NotAfter: now.Add(time.Hour), DNSNames: []string{dns},
|
||||
BasicConstraintsValid: true, KeyUsage: x509.KeyUsageDigitalSignature,
|
||||
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth, x509.ExtKeyUsageClientAuth},
|
||||
}
|
||||
der, err := x509.CreateCertificate(rand.Reader, template, root, &key.PublicKey, caKey)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
parsed, err := x509.ParseCertificate(der)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
privateDER, err := x509.MarshalPKCS8PrivateKey(key)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}), pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: privateDER}), parsed
|
||||
}
|
||||
agentCert, agentKey, _ = issue(2, "agent.local")
|
||||
dispatcherCert, dispatcherKey, dispatcherLeaf = issue(3, "dispatcher.local")
|
||||
return
|
||||
}
|
||||
@@ -52,7 +52,7 @@ func newRootCommand() *cobra.Command {
|
||||
SilenceUsage: true,
|
||||
SilenceErrors: true,
|
||||
}
|
||||
root.AddCommand(newAgentCommand(), newDispatcherCommand())
|
||||
root.AddCommand(newCurrentAgentCommand(), newDispatcherCommand())
|
||||
return root
|
||||
}
|
||||
|
||||
|
||||
@@ -19,12 +19,11 @@ func TestRootHasExplicitRoles(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentRejectsInvalidMediaPortEnvironment(t *testing.T) {
|
||||
t.Setenv("AGENT_CALL_MEDIA_PORT", "invalid")
|
||||
func TestRootAgentRejectsOldOneShotFlag(t *testing.T) {
|
||||
root := newRootCommand()
|
||||
root.SetArgs([]string{"agent"})
|
||||
if err := root.Execute(); err == nil || !strings.Contains(err.Error(), "AGENT_CALL_MEDIA_PORT") {
|
||||
t.Fatalf("expected invalid port error, got %v", err)
|
||||
root.SetArgs([]string{"agent", "--call-once"})
|
||||
if err := root.Execute(); err == nil || !strings.Contains(err.Error(), "unknown flag") {
|
||||
t.Fatalf("obsolete one-shot entry was admitted: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user