Assemble isolated Dispatcher runtime command
This commit is contained in:
@@ -0,0 +1,170 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net"
|
||||
"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/oss"
|
||||
"git.ipao.vip/rogee/go-sip/internal/rpc"
|
||||
"git.ipao.vip/rogee/go-sip/internal/store"
|
||||
"github.com/google/uuid"
|
||||
"github.com/spf13/cobra"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials"
|
||||
)
|
||||
|
||||
func newCurrentDispatcherCommand() *cobra.Command {
|
||||
var mode string
|
||||
command := &cobra.Command{
|
||||
Use: "dispatcher", Short: "Run the isolated Dispatcher", Args: cobra.NoArgs,
|
||||
SilenceUsage: true,
|
||||
RunE: func(cmd *cobra.Command, _ []string) error {
|
||||
return runCurrentDispatcher(cmd.Context(), mode)
|
||||
},
|
||||
}
|
||||
command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock mode only")
|
||||
return command
|
||||
}
|
||||
|
||||
// runCurrentDispatcher does not publish or consume until the predeclared MQ
|
||||
// topology, Agent session and approved SIP load have all been verified.
|
||||
func runCurrentDispatcher(ctx context.Context, mode string) (result error) {
|
||||
settings, err := config.LoadDispatcherRuntimeEnvironment(mode)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
agentEndpoint, err := config.LoadCurrentMockAgentEndpoint(settings.AgentEndpointsFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
ossConfiguration, err := config.LoadCurrentOSSConfig(settings.OSSConfigFile, settings.DispatcherID)
|
||||
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("Dispatcher mTLS listener certificate is invalid")
|
||||
}
|
||||
clientTLS, err := rpc.NewClientTLSConfig(ca, cert, key, agentEndpoint.ServerName)
|
||||
if err != nil {
|
||||
return errors.New("Dispatcher mTLS Agent certificate configuration is invalid")
|
||||
}
|
||||
httpClient := localMockHTTPClient()
|
||||
httpClient.Timeout = 10 * time.Second
|
||||
_, reader, err := dispatcherConfigurationClient(mode, httpClient)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
ossClient, err := oss.NewClient(ossConfiguration)
|
||||
if err != nil {
|
||||
return fmt.Errorf("initialize bounded OSS grant signer: %w", err)
|
||||
}
|
||||
database, err := store.OpenCurrent(settings.SQLitePath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { result = errors.Join(result, database.Close()) }()
|
||||
if err := database.CloseAdmission(settings.DispatcherID); err != nil {
|
||||
return fmt.Errorf("close task admission before connecting external adapters: %w", err)
|
||||
}
|
||||
broker, err := mq.OpenCurrent(settings.RabbitMQURL, settings.DispatcherID, 1)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { result = errors.Join(result, broker.Close()) }()
|
||||
connection, err := grpc.NewClient(agentEndpoint.Address, grpc.WithTransportCredentials(credentials.NewTLS(clientTLS)))
|
||||
if err != nil {
|
||||
return errors.New("Dispatcher cannot create verified Agent connection")
|
||||
}
|
||||
defer func() { result = errors.Join(result, connection.Close()) }()
|
||||
agentClient := agentpb.NewAgentControlServiceClient(connection)
|
||||
coordinator := dispatcher.NewAgentCoordinator(time.Now)
|
||||
if err := coordinator.Register(agentEndpoint.AgentID, agentClient); err != nil {
|
||||
return err
|
||||
}
|
||||
probeCtx, cancelProbe := context.WithTimeout(ctx, 5*time.Second)
|
||||
status, err := coordinator.Probe(probeCtx, agentEndpoint.AgentID, agentEndpoint.CellID)
|
||||
if err != nil {
|
||||
cancelProbe()
|
||||
return fmt.Errorf("probe assigned Agent before admission: %w", err)
|
||||
}
|
||||
epoch, err := uuid.NewRandom()
|
||||
if err != nil {
|
||||
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 {
|
||||
cancelProbe()
|
||||
return fmt.Errorf("activate assigned Agent before admission: %w", err)
|
||||
}
|
||||
cancelProbe()
|
||||
originator := &dispatcher.ApprovedOriginator{
|
||||
DispatcherID: settings.DispatcherID,
|
||||
Client: agentClient,
|
||||
Meta: func(callCtx context.Context) (*agentpb.RequestMeta, error) {
|
||||
return coordinator.ApprovedMeta(callCtx, agentEndpoint.AgentID)
|
||||
},
|
||||
}
|
||||
worker := &dispatcher.CurrentRuntime{
|
||||
Broker: broker,
|
||||
Bootstrap: dispatcher.CurrentBootstrap{
|
||||
DispatcherID: settings.DispatcherID, Client: reader, Store: database, VerifySIP: originator.VerifySIP,
|
||||
},
|
||||
Execute: dispatcher.CurrentExecuteController{
|
||||
DispatcherID: settings.DispatcherID, Store: database, Originator: originator, Publisher: broker, Now: time.Now,
|
||||
},
|
||||
Control: dispatcher.CurrentControlController{
|
||||
DispatcherID: settings.DispatcherID, Store: database, Client: reader, Agent: originator, VerifySIP: originator.VerifySIP, Now: time.Now,
|
||||
},
|
||||
PollInterval: time.Second, DiscoveryInterval: time.Minute, Logger: slog.Default(),
|
||||
}
|
||||
listener, err := net.Listen("tcp", settings.Listen)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Dispatcher cannot listen on configured Mock address: %w", err)
|
||||
}
|
||||
server := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS)))
|
||||
agentpb.RegisterAgentControlServiceServer(server, &rpc.RecordingServer{
|
||||
DispatcherID: settings.DispatcherID, Store: database, OSS: ossClient,
|
||||
TrustedFingerprints: settings.PeerFingerprints, AuthorizeSession: coordinator.AuthorizeInboundMeta,
|
||||
})
|
||||
serveCtx, cancelServe := context.WithCancel(ctx)
|
||||
defer cancelServe()
|
||||
serverErrors := make(chan error, 1)
|
||||
go func() {
|
||||
serverErrors <- server.Serve(listener)
|
||||
cancelServe()
|
||||
}()
|
||||
runtimeErr := worker.Serve(serveCtx)
|
||||
server.GracefulStop()
|
||||
serverErr := <-serverErrors
|
||||
if serverErr != nil && !errors.Is(serverErr, grpc.ErrServerStopped) {
|
||||
return errors.Join(runtimeErr, fmt.Errorf("Dispatcher recording listener stopped: %w", serverErr))
|
||||
}
|
||||
if runtimeErr != nil {
|
||||
return runtimeErr
|
||||
}
|
||||
if ctx.Err() == nil {
|
||||
return errors.New("Dispatcher stopped without process shutdown")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,83 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestCurrentDispatcherCommandRejectsOldFlagsAndModesBeforeResources(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
args []string
|
||||
want string
|
||||
}{
|
||||
{"real", []string{"--mode", "real"}, "Mock"},
|
||||
{"mixed", []string{"--mode", "mixed"}, "Mock"},
|
||||
{"old config", []string{"--config", "/nonexistent"}, "unknown flag"},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
database := filepath.Join(t.TempDir(), "dispatcher.sqlite")
|
||||
t.Setenv("DISPATCHER_SQLITE_PATH", database)
|
||||
command := newCurrentDispatcherCommand()
|
||||
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("unapproved Dispatcher entry was admitted: %v", err)
|
||||
}
|
||||
if _, err := os.Stat(database); !os.IsNotExist(err) {
|
||||
t.Fatalf("rejected entry opened business state: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentDispatcherCommandValidatesDeploymentBeforeOpeningSQLite(t *testing.T) {
|
||||
root := t.TempDir()
|
||||
endpointFile := filepath.Join(root, "agent-endpoints.json")
|
||||
ossFile := filepath.Join(root, "oss.json")
|
||||
database := filepath.Join(root, "dispatcher.sqlite")
|
||||
for name, value := range map[string]string{
|
||||
"DISPATCHER_ID": "c046b893-8628-4589-ae50-619d049248a6",
|
||||
"DISPATCHER_SECRET_KEY": "isolated-placeholder",
|
||||
"SAAS_BASE_URL": "http://127.0.0.1:18770",
|
||||
"RABBITMQ_URL": "amqp://127.0.0.1:18771",
|
||||
"DISPATCHER_SQLITE_PATH": database,
|
||||
"DISPATCHER_GRPC_LISTEN": "127.0.0.1:0",
|
||||
"DISPATCHER_AGENT_ENDPOINTS_FILE": endpointFile,
|
||||
"DISPATCHER_OSS_CONFIG_FILE": ossFile,
|
||||
"MTLS_CA_FILE": filepath.Join(root, "ca.pem"),
|
||||
"MTLS_CERT_FILE": filepath.Join(root, "dispatcher.pem"),
|
||||
"MTLS_KEY_FILE": filepath.Join(root, "dispatcher.key"),
|
||||
"MTLS_PEER_CERT_FINGERPRINTS": strings.Repeat("a", 64),
|
||||
} {
|
||||
t.Setenv(name, value)
|
||||
}
|
||||
if err := os.WriteFile(endpointFile, []byte(`[{"agent_id":"agent-mock","cell_id":"cell-mock","address":"saas.example.invalid:443","server_name":"agent.local"}]`), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(ossFile, []byte(`{"dispatcher_id":"22222222-2222-4222-8222-222222222222","oss":{}}`), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
check := func(want string) {
|
||||
t.Helper()
|
||||
command := newCurrentDispatcherCommand()
|
||||
command.SetOut(io.Discard)
|
||||
command.SetErr(io.Discard)
|
||||
command.SetArgs([]string{"--mode", "mock"})
|
||||
if err := command.Execute(); err == nil || !strings.Contains(err.Error(), want) || strings.Contains(err.Error(), "saas.example.invalid") {
|
||||
t.Fatalf("unapproved Dispatcher deployment was admitted or echoed: %v", err)
|
||||
}
|
||||
if _, err := os.Stat(database); !os.IsNotExist(err) {
|
||||
t.Fatalf("rejected deployment opened SQLite: %v", err)
|
||||
}
|
||||
}
|
||||
check("local")
|
||||
if err := os.WriteFile(endpointFile, []byte(`[{"agent_id":"agent-mock","cell_id":"cell-mock","address":"127.0.0.1:19090","server_name":"agent.local"}]`), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
check("not assigned")
|
||||
}
|
||||
@@ -35,6 +35,7 @@
|
||||
- 新增无实现代次字段的 Dispatcher 环境预检:必须显式给出规范 UUID v4 身份、只读 HTTP 地址及密钥、RabbitMQ 地址和 SQLite 路径;当前 Mock 仅接受本机 HTTP/MQ 目标,mixed/real 在任何资源操作前拒绝。预检不打开数据库或网络,缺失值不继承旧默认配置,错误不打印凭据。`go test ./internal/config -run '^TestLoadDispatcherEnvironment' -count=1` 通过;此预检尚未接入主 CLI,不能当作 P02/P03 完成。
|
||||
- Dispatcher 另有纯运行环境预检:唯一 Agent 端点清单和 OSS 配置文件须提供路径,本机 gRPC 监听及双向 TLS 文件/Agent 证书指纹须显式给出;缺失、非法指纹或非本机监听明确拒绝且不回显值。该检查不读取配置文件、不连接 Agent、不打开 SQLite/MQ;`go test ./internal/config -run '^TestLoadDispatcherRuntimeEnvironment' -count=1` 通过。当前仍未由 `dispatcher` 主命令调用,不能当作 Agent 加载或 OSS 签发已验收。
|
||||
- D 的部署文件读取已独立于旧代次 Schema:Agent 清单仅接受**一个本机端点**;OSS 配置必须匹配本 D 身份、指定本机 HTTPS 目标、15 分钟授权与显式资产上限,凭据只经受控环境变量引用。超大文件、重复/未知字段、多个 Agent、非本机目标、缺失凭据和错误 D 归属均拒绝;JSON 唯一键及凭据引用基础函数已移出旧配置文件。`go test ./internal/config -count=1` 通过。此处仅验证部署数据,不代表 Agent 实际加载、官方 SDK 已签发目标,亦未接入主 `dispatcher` 命令。
|
||||
- 独立的隔离 D 命令函数现已组装严格部署预检、当前 HTTP 读取、`store.OpenCurrent`、被动核验 SaaS 预建拓扑的 `mq.OpenCurrent`、D 身份绑定的 Agent 会话、批准执行/控制、实际加载版本检查、`CurrentRuntime` 和录音双向 TLS RPC;启动开库后先持久关闭新准入,不清理旧状态或擅自建队。负例测试确认 mixed/real、旧 `--config` 与非本机 Agent/错误 OSS 归属在打开 SQLite 前拒绝。**根 `dispatcher` 命令尚未切换;该函数尚未用实际本机 Agent、RabbitMQ 和 HTTP 完整启动,不构成主程序通过或外部验收。**
|
||||
|
||||
- D 的 Agent 会话激活现在必须显式携带规范 UUID v4 的本 D 身份,并写入 Agent 授权绑定;预配置了对应 D 身份的隔离 Agent 会拒绝错误归属,不能再靠 D 自报空身份放行。`go test ./internal/dispatcher -count=1` 通过;主 `dispatcher` 命令尚未使用这条会话链路,不代表本机 D↔A 执行已验收。
|
||||
- D 的已授权调用元数据现在只从该 Agent **最新且未过期**的会话产生,每次生成不同操作身份;未注册、会话缺失或过期均明确拒绝。`TestAgentCoordinatorProbesBeforeActivation` 覆盖这些边界,不能代替实际 D 主进程的执行与录音联测。
|
||||
|
||||
Reference in New Issue
Block a user