refactor: give business commands and deployment config domain names

This commit is contained in:
2026-09-30 17:33:58 +08:00
parent 04b66c4c4c
commit 1647aefd5b
22 changed files with 135 additions and 55 deletions
@@ -18,7 +18,7 @@ import (
"google.golang.org/grpc/credentials"
)
func newCurrentAgentCommand() *cobra.Command {
func newAgentCommand() *cobra.Command {
var mode string
command := &cobra.Command{
Use: "agent", Short: "Run the isolated Agent", Args: cobra.NoArgs,
@@ -61,7 +61,7 @@ func newCurrentAgentCommand() *cobra.Command {
return errors.New("Agent cannot create pinned Dispatcher connection")
}
defer connection.Close()
handler, err := newCurrentAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection))
handler, err := newAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection))
if err != nil {
return err
}
@@ -18,7 +18,7 @@ import (
"google.golang.org/grpc/status"
)
func TestCurrentAgentCommandRejectsRealModesAndObsoleteCallEntry(t *testing.T) {
func TestAgentCommandRejectsRealModesAndObsoleteCallEntry(t *testing.T) {
for _, tc := range []struct {
name string
args []string
@@ -33,7 +33,7 @@ func TestCurrentAgentCommandRejectsRealModesAndObsoleteCallEntry(t *testing.T) {
state := filepath.Join(t.TempDir(), "session.json")
t.Setenv("AGENT_ID", "")
t.Setenv("AGENT_SESSION_PATH", state)
command := newCurrentAgentCommand()
command := newAgentCommand()
command.SetOut(io.Discard)
command.SetErr(io.Discard)
command.SetArgs(tc.args)
@@ -73,7 +73,7 @@ func TestLoadMockAppliedSIPAcceptsOnlyExplicitLocalFacts(t *testing.T) {
}
}
func TestCurrentAgentCommandServesPinnedMutualTLSOnly(t *testing.T) {
func TestAgentCommandServesPinnedMutualTLSOnly(t *testing.T) {
ca, agentCert, agentKey, dispatcherCert, dispatcherKey, dispatcherLeaf := localCommandCertificates(t)
settings, _ := currentAgentSetupFixture(t)
for path, body := range map[string][]byte{
@@ -112,7 +112,7 @@ func TestCurrentAgentCommandServesPinnedMutualTLSOnly(t *testing.T) {
t.Setenv(name, value)
}
process, cancel := context.WithCancel(context.Background())
command := newCurrentAgentCommand()
command := newAgentCommand()
command.SetOut(io.Discard)
command.SetErr(io.Discard)
command.SetContext(process)
@@ -26,38 +26,38 @@ import (
"github.com/google/uuid"
)
// newCurrentAgentServer binds the authenticated Agent session, task controls
// newAgentServer binds the authenticated Agent session, task controls
// and per-call Mock recording delivery. Serving the returned server requires
// a separately verified mutual-TLS listener and a pinned local D connection.
func newCurrentAgentServer(ctx context.Context, settings config.AgentEnvironment, scenario approvedMockScenario, appliedSIP map[string]int64, dispatcher agentpb.AgentControlServiceClient) (*rpc.Server, error) {
func newAgentServer(ctx context.Context, settings config.AgentEnvironment, scenario approvedMockScenario, appliedSIP map[string]int64, dispatcher agentpb.AgentControlServiceClient) (*rpc.Server, error) {
if ctx == nil || ctx.Err() != nil || settings.AgentID == "" || settings.CellID == "" || settings.SessionPath == "" ||
settings.RecoveryRoot == "" || len(settings.PeerFingerprints) == 0 || len(appliedSIP) == 0 ||
scenario.MaxWAVBytes <= 44 || len(scenario.Script.Turns) == 0 || strings.TrimSpace(scenario.ReasonMessage) == "" {
return nil, errors.New("current Agent requires an active process, explicit Mock media and deployment identity")
return nil, errors.New("Mock Agent requires an active process, explicit Mock media and deployment identity")
}
if tenant.ValidateDispatcherID(settings.DispatcherID) != nil {
return nil, errors.New("current Agent requires an approved Dispatcher UUID v4")
return nil, errors.New("Mock Agent requires an approved Dispatcher UUID v4")
}
if dispatcher == nil {
return nil, errors.New("current Agent requires a pinned Dispatcher transport")
return nil, errors.New("Mock Agent requires a pinned Dispatcher transport")
}
value := reflect.ValueOf(dispatcher)
switch value.Kind() {
case reflect.Chan, reflect.Func, reflect.Interface, reflect.Map, reflect.Pointer, reflect.Slice:
if value.IsNil() {
return nil, errors.New("current Agent requires a pinned Dispatcher transport")
return nil, errors.New("Mock Agent requires a pinned Dispatcher transport")
}
}
root, err := os.Stat(settings.RecoveryRoot)
if err != nil || !root.IsDir() || root.Mode().Perm() != 0700 {
return nil, errors.New("current Agent requires an existing private 0700 recovery directory")
return nil, errors.New("Mock Agent requires an existing private 0700 recovery directory")
}
if err := rejectLegacyAgentSpool(settings.RecoveryRoot); err != nil {
return nil, err
}
for trunk, revision := range appliedSIP {
if strings.TrimSpace(trunk) == "" || revision <= 0 {
return nil, errors.New("current Agent requires explicit applied Mock SIP revisions")
return nil, errors.New("Mock Agent requires explicit applied Mock SIP revisions")
}
}
loaded := maps.Clone(appliedSIP)
@@ -122,7 +122,7 @@ func newCurrentAgentServer(ctx context.Context, settings config.AgentEnvironment
}
// Old spool files may contain unreported execution or upload outcomes. Never
// start the current Agent on the same root without an explicit disposition.
// start another Agent on the same root without an explicit disposition.
func rejectLegacyAgentSpool(root string) error {
entries, err := os.ReadDir(root)
if err != nil {
@@ -47,9 +47,9 @@ func currentAgentSetupFixture(t *testing.T) (config.AgentEnvironment, approvedMo
return settings, scenario
}
func TestNewCurrentAgentServerBindsMockCallsToOneSessionAndRecoveryRoot(t *testing.T) {
func TestNewAgentServerBindsMockCallsToOneSessionAndRecoveryRoot(t *testing.T) {
settings, scenario := currentAgentSetupFixture(t)
server, err := newCurrentAgentServer(context.Background(), settings, scenario, map[string]int64{"trunk-mock": 8}, &isolatedAgentRecordingClient{})
server, err := newAgentServer(context.Background(), settings, scenario, map[string]int64{"trunk-mock": 8}, &isolatedAgentRecordingClient{})
if err != nil || server == nil {
t.Fatalf("valid isolated Agent could not be assembled: %v", err)
}
@@ -58,7 +58,7 @@ func TestNewCurrentAgentServerBindsMockCallsToOneSessionAndRecoveryRoot(t *testi
}
}
func TestNewCurrentAgentServerRefusesLegacySpoolWithoutChangingIt(t *testing.T) {
func TestNewAgentServerRefusesLegacySpoolWithoutChangingIt(t *testing.T) {
for _, tc := range []struct {
name string
marker string
@@ -83,7 +83,7 @@ func TestNewCurrentAgentServerRefusesLegacySpoolWithoutChangingIt(t *testing.T)
t.Fatal(err)
}
}
server, err := newCurrentAgentServer(context.Background(), settings, scenario, map[string]int64{"trunk-mock": 8}, &isolatedAgentRecordingClient{})
server, err := newAgentServer(context.Background(), settings, scenario, map[string]int64{"trunk-mock": 8}, &isolatedAgentRecordingClient{})
if err == nil || server != nil {
t.Fatalf("legacy spool was admitted: %v", err)
}
@@ -97,7 +97,7 @@ func TestNewCurrentAgentServerRefusesLegacySpoolWithoutChangingIt(t *testing.T)
}
}
func TestNewCurrentAgentServerRefusesUnsafeAdaptersBeforeResources(t *testing.T) {
func TestNewAgentServerRefusesUnsafeAdaptersBeforeResources(t *testing.T) {
settings, scenario := currentAgentSetupFixture(t)
for _, tc := range []struct {
name string
@@ -134,7 +134,7 @@ func TestNewCurrentAgentServerRefusesUnsafeAdaptersBeforeResources(t *testing.T)
loaded := map[string]int64{"trunk-mock": 8}
var client agentpb.AgentControlServiceClient = &isolatedAgentRecordingClient{}
tc.change(&cfg, &media, &loaded, &client)
if server, err := newCurrentAgentServer(context.Background(), cfg, media, loaded, client); err == nil || server != nil {
if server, err := newAgentServer(context.Background(), cfg, media, loaded, client); err == nil || server != nil {
t.Fatalf("unsafe current Agent assembly was accepted: %v", err)
}
if _, err := os.Stat(settings.SessionPath); !os.IsNotExist(err) {
@@ -21,31 +21,31 @@ import (
"google.golang.org/grpc/credentials"
)
func newCurrentDispatcherCommand() *cobra.Command {
func newDispatcherCommand() *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)
return runDispatcher(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
// runDispatcher 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) {
func runDispatcher(ctx context.Context, mode string) (result error) {
settings, err := config.LoadDispatcherRuntimeEnvironment(mode)
if err != nil {
return err
}
agentEndpoint, err := config.LoadCurrentMockAgentEndpoint(settings.AgentEndpointsFile)
agentEndpoint, err := config.LoadMockAgentEndpoint(settings.AgentEndpointsFile)
if err != nil {
return err
}
ossConfiguration, err := config.LoadCurrentOSSConfig(settings.OSSConfigFile, settings.DispatcherID)
ossConfiguration, err := config.LoadOSSConfig(settings.OSSConfigFile, settings.DispatcherID)
if err != nil {
return err
}
@@ -8,7 +8,7 @@ import (
"testing"
)
func TestCurrentDispatcherCommandRejectsOldFlagsAndModesBeforeResources(t *testing.T) {
func TestDispatcherCommandRejectsOldFlagsAndModesBeforeResources(t *testing.T) {
for _, tc := range []struct {
name string
args []string
@@ -21,7 +21,7 @@ func TestCurrentDispatcherCommandRejectsOldFlagsAndModesBeforeResources(t *testi
t.Run(tc.name, func(t *testing.T) {
database := filepath.Join(t.TempDir(), "dispatcher.sqlite")
t.Setenv("DISPATCHER_SQLITE_PATH", database)
command := newCurrentDispatcherCommand()
command := newDispatcherCommand()
command.SetOut(io.Discard)
command.SetErr(io.Discard)
command.SetArgs(tc.args)
@@ -35,7 +35,7 @@ func TestCurrentDispatcherCommandRejectsOldFlagsAndModesBeforeResources(t *testi
}
}
func TestCurrentDispatcherCommandValidatesDeploymentBeforeOpeningSQLite(t *testing.T) {
func TestDispatcherCommandValidatesDeploymentBeforeOpeningSQLite(t *testing.T) {
root := t.TempDir()
endpointFile := filepath.Join(root, "agent-endpoints.json")
ossFile := filepath.Join(root, "oss.json")
@@ -64,7 +64,7 @@ func TestCurrentDispatcherCommandValidatesDeploymentBeforeOpeningSQLite(t *testi
}
check := func(want string) {
t.Helper()
command := newCurrentDispatcherCommand()
command := newDispatcherCommand()
command.SetOut(io.Discard)
command.SetErr(io.Discard)
command.SetArgs([]string{"--mode", "mock"})
+1 -1
View File
@@ -9,7 +9,7 @@ import (
)
func TestDispatcherHasNoBusinessHTTPFlags(t *testing.T) {
command := newCurrentDispatcherCommand()
command := newDispatcherCommand()
for _, name := range []string{"control-listen", "control-token", "config", "db"} {
if command.Flags().Lookup(name) != nil {
t.Fatalf("obsolete Dispatcher flag remains: --%s", name)
@@ -32,7 +32,7 @@ import (
"google.golang.org/grpc/credentials"
)
func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) {
func TestDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) {
brokerURL, adminURL := os.Getenv("RABBITMQ_URL"), os.Getenv("RABBITMQ_PROVISIONER_URL")
if brokerURL == "" || adminURL == "" {
t.Skip("requires the isolated RabbitMQ mock provisioner")
@@ -165,7 +165,7 @@ func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T)
defer agentConnection.Close()
agentContext, stopAgent := context.WithCancel(context.Background())
defer stopAgent()
agentServer, err := newCurrentAgentServer(agentContext, settings, scenario, map[string]int64{"trunk-mock": 8}, agentpb.NewAgentControlServiceClient(agentConnection))
agentServer, err := newAgentServer(agentContext, settings, scenario, map[string]int64{"trunk-mock": 8}, agentpb.NewAgentControlServiceClient(agentConnection))
if err != nil {
t.Fatal(err)
}
+1 -1
View File
@@ -21,6 +21,6 @@ func newRootCommand() *cobra.Command {
SilenceUsage: true,
SilenceErrors: true,
}
root.AddCommand(newCurrentAgentCommand(), newCurrentDispatcherCommand())
root.AddCommand(newAgentCommand(), newDispatcherCommand())
return root
}
+39
View File
@@ -0,0 +1,39 @@
package main
import (
"go/ast"
"go/parser"
"go/token"
"os"
"strings"
"testing"
)
func TestBusinessCommandsUseDomainNames(t *testing.T) {
files, err := os.ReadDir(".")
if err != nil {
t.Fatal(err)
}
for _, file := range files {
name := file.Name()
if file.IsDir() || !strings.HasSuffix(name, ".go") {
continue
}
if strings.HasPrefix(name, "current_") {
t.Errorf("%s retains an implementation-generation filename", name)
}
if strings.HasSuffix(name, "_test.go") {
continue
}
parsed, err := parser.ParseFile(token.NewFileSet(), name, nil, 0)
if err != nil {
t.Fatal(err)
}
ast.Inspect(parsed, func(node ast.Node) bool {
if declaration, ok := node.(*ast.FuncDecl); ok && strings.Contains(declaration.Name.Name, "Current") {
t.Errorf("%s retains an implementation-generation function %s", name, declaration.Name.Name)
}
return true
})
}
}