diff --git a/cmd/sip-go-agent/main.go b/cmd/sip-go-agent/main.go index 12de036..91ee6f5 100644 --- a/cmd/sip-go-agent/main.go +++ b/cmd/sip-go-agent/main.go @@ -57,12 +57,15 @@ func newRootCommand() *cobra.Command { } func newAgentCommand() *cobra.Command { - cfg := config.FromEnv() + cfg, configErr := config.FromEnv() var realCall bool cmd := &cobra.Command{ Use: "agent", Short: "run the file-backed Agent process", RunE: func(_ *cobra.Command, _ []string) error { + if configErr != nil { + return configErr + } if realCall { return runCallOnce(cfg) } @@ -469,13 +472,16 @@ func openDispatcherStore(cfg config.Config) (*store.Store, error) { } func newDispatcherCommand() *cobra.Command { - cfg := config.FromEnv() + cfg, configErr := config.FromEnv() var configFile string var once bool cmd := &cobra.Command{ Use: "dispatcher", Short: "run the single-active Dispatcher process", RunE: func(cmd *cobra.Command, _ []string) error { + if configErr != nil { + return configErr + } if err := cfg.LoadDispatcherFile(configFile); err != nil { return err } diff --git a/cmd/sip-go-agent/main_test.go b/cmd/sip-go-agent/main_test.go index 7b5e724..2021961 100644 --- a/cmd/sip-go-agent/main_test.go +++ b/cmd/sip-go-agent/main_test.go @@ -2,6 +2,7 @@ package main import ( "bytes" + "strings" "testing" "git.ipao.vip/rogee/go-sip/contracts" @@ -18,6 +19,15 @@ func TestRootHasExplicitRoles(t *testing.T) { } } +func TestAgentRejectsInvalidMediaPortEnvironment(t *testing.T) { + t.Setenv("AGENT_CALL_MEDIA_PORT", "invalid") + 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) + } +} + func TestWriteResultIsJSON(t *testing.T) { var b bytes.Buffer data, err := marshalResultForTest(map[string]string{"status": "ok"}) diff --git a/cmd/sip-go-agent/upload_retry_command.go b/cmd/sip-go-agent/upload_retry_command.go index 94a19fa..0f72aea 100644 --- a/cmd/sip-go-agent/upload_retry_command.go +++ b/cmd/sip-go-agent/upload_retry_command.go @@ -11,12 +11,15 @@ import ( ) func newUploadRetryCommand() *cobra.Command { - cfg := config.FromEnv() + cfg, configErr := config.FromEnv() var uploadID, requestID, path string cmd := &cobra.Command{ Use: "upload-retry", Short: "explicitly request one new grant for an unsuccessful upload; never dial", Args: cobra.NoArgs, RunE: func(cmd *cobra.Command, _ []string) (returnErr error) { + if configErr != nil { + return configErr + } if uploadID == "" || requestID == "" || path == "" { return errors.New("--upload-id, --request-id and --file are required") } diff --git a/internal/ai/snapshot.go b/internal/ai/snapshot.go index 9b42de6..6d7db77 100644 --- a/internal/ai/snapshot.go +++ b/internal/ai/snapshot.go @@ -9,7 +9,6 @@ import ( "encoding/json" "errors" "fmt" - "sync" "unicode/utf8" "github.com/cyberphone/json-canonicalization/go/src/webpki.org/jsoncanonicalizer" @@ -78,47 +77,4 @@ func ValidateForMode(raw []byte, mode Mode) (Snapshot, error) { return snapshot, nil } -func EnsureSameVersion(previous, next Snapshot) error { - if previous.AgentVersionID == "" || next.AgentVersionID == "" || previous.AgentVersionID != next.AgentVersionID { - return errors.New("agent version identity changed") - } - if previous.Digest != next.Digest { - return errors.New("immutable agent version content changed") - } - return nil -} - -type Cache struct { - mu sync.RWMutex - items map[string]Snapshot -} - -func NewCache() *Cache { return &Cache{items: make(map[string]Snapshot)} } - -func (c *Cache) Put(tenantKey string, snapshot Snapshot) error { - if err := contract.ValidateTenantKey(tenantKey); err != nil { - return err - } - if snapshot.AgentVersionID == "" || snapshot.Digest == "" { - return errors.New("snapshot identity is required") - } - c.mu.Lock() - defer c.mu.Unlock() - key := tenantKey + "\x00" + snapshot.AgentVersionID - if old, ok := c.items[key]; ok { - if err := EnsureSameVersion(old, snapshot); err != nil { - return err - } - } - c.items[key] = snapshot - return nil -} - -func (c *Cache) Get(tenantKey, versionID string) (Snapshot, bool) { - c.mu.RLock() - defer c.mu.RUnlock() - s, ok := c.items[tenantKey+"\x00"+versionID] - return s, ok -} - func ContractSource() string { return contracts.SourceCommit } diff --git a/internal/ai/snapshot_test.go b/internal/ai/snapshot_test.go index 962e105..5f61977 100644 --- a/internal/ai/snapshot_test.go +++ b/internal/ai/snapshot_test.go @@ -43,23 +43,3 @@ func TestASROnlyRejectsFullAIFields(t *testing.T) { t.Fatal("expected ASR-only schema rejection") } } - -func TestCacheRejectsChangedImmutableContent(t *testing.T) { - raw, err := contracts.Read("examples/agent-version.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := Validate(raw) - if err != nil { - t.Fatal(err) - } - cache := NewCache() - if err := cache.Put("tenant-demo-key", snapshot); err != nil { - t.Fatal(err) - } - changed := snapshot - changed.Digest = "different" - if err := cache.Put("tenant-demo-key", changed); err == nil { - t.Fatal("expected immutable content rejection") - } -} diff --git a/internal/config/config.go b/internal/config/config.go index 7bd0907..322b057 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -61,7 +61,11 @@ type Config struct { CallExecutionID string } -func FromEnv() Config { +func FromEnv() (Config, error) { + mediaPort, err := envInt("AGENT_CALL_MEDIA_PORT", 12000) + if err != nil { + return Config{}, err + } return Config{ Mode: envOr("SIP_GO_AGENT_MODE", "mock"), DBPath: envOr("DISPATCHER_DB", "dispatcher.db"), @@ -97,7 +101,7 @@ func FromEnv() Config { CallTrunkID: os.Getenv("AGENT_CALL_TRUNK_ID"), CallCallerID: os.Getenv("AGENT_CALL_CALLER_ID"), CallMediaBind: envOr("AGENT_CALL_MEDIA_BIND", "127.0.0.1"), - CallMediaPort: envInt("AGENT_CALL_MEDIA_PORT", 12000), + CallMediaPort: mediaPort, CallRecordingDirectory: envOr("AGENT_CALL_RECORDING_DIR", "./recordings"), CallAISnapshotPath: os.Getenv("AGENT_CALL_AI_SNAPSHOT"), CallTenantID: os.Getenv("AGENT_CALL_TENANT_ID"), @@ -105,7 +109,7 @@ func FromEnv() Config { CallTaskID: os.Getenv("AGENT_CALL_TASK_ID"), CallTaskItemID: os.Getenv("AGENT_CALL_TASK_ITEM_ID"), CallExecutionID: os.Getenv("AGENT_CALL_EXECUTION_ID"), - } + }, nil } func (c Config) Validate(role string) error { @@ -191,29 +195,14 @@ func envOr(key, fallback string) string { return fallback } -func envOrSecret(valueKey, fileKey string) string { - if value := strings.TrimSpace(os.Getenv(valueKey)); value != "" { - return value - } - path := strings.TrimSpace(os.Getenv(fileKey)) - if path == "" { - return "" - } - data, err := os.ReadFile(path) - if err != nil { - return "" - } - return strings.TrimSpace(string(data)) -} - -func envInt(key string, fallback int) int { +func envInt(key string, fallback int) (int, error) { value := strings.TrimSpace(os.Getenv(key)) if value == "" { - return fallback + return fallback, nil } parsed, err := strconv.Atoi(value) if err != nil { - return fallback + return 0, fmt.Errorf("invalid %s: %w", key, err) } - return parsed + return parsed, nil } diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 3e21cc8..dfceb8f 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -1,15 +1,38 @@ package config -import "testing" +import ( + "strings" + "testing" +) func TestConfigRejectsRealDispatcherWithoutBroker(t *testing.T) { - c := FromEnv() + c, err := FromEnv() + if err != nil { + t.Fatal(err) + } c.Mode, c.RabbitURL, c.DBPath = "real", "", ":memory:" if err := c.Validate("dispatcher"); err == nil { t.Fatal("expected real-mode broker requirement") } } +func TestFromEnvRejectsInvalidMediaPort(t *testing.T) { + t.Setenv("AGENT_CALL_MEDIA_PORT", "not-a-port") + if _, err := FromEnv(); err == nil || !strings.Contains(err.Error(), "AGENT_CALL_MEDIA_PORT") { + t.Fatalf("expected explicit media port error, got %v", err) + } + t.Setenv("AGENT_CALL_MEDIA_PORT", "43210") + cfg, err := FromEnv() + if err != nil || cfg.CallMediaPort != 43210 { + t.Fatalf("expected configured port, got %d, %v", cfg.CallMediaPort, err) + } + t.Setenv("AGENT_CALL_MEDIA_PORT", "") + cfg, err = FromEnv() + if err != nil || cfg.CallMediaPort != 12000 { + t.Fatalf("expected default port, got %d, %v", cfg.CallMediaPort, err) + } +} + func TestConfigAcceptsMockAgent(t *testing.T) { c := Config{Mode: "mock", SpoolRoot: t.TempDir()} if err := c.Validate("agent"); err != nil { @@ -22,7 +45,10 @@ func TestConfigLoadsUnifiedDispatcherGRPCSettings(t *testing.T) { t.Setenv("DISPATCHER_GRPC_ENDPOINT", "dispatcher.test:19443") t.Setenv("DISPATCHER_GRPC_SERVER_NAME", "dispatcher.test") t.Setenv("DISPATCHER_ALLOWED_AGENT_IDS", "agent-cell-a,agent-cell-b") - c := FromEnv() + c, err := FromEnv() + if err != nil { + t.Fatal(err) + } if c.DispatcherGRPCListen != "127.0.0.1:19443" || c.DispatcherGRPCEndpoint != "dispatcher.test:19443" || c.DispatcherGRPCServerName != "dispatcher.test" || c.DispatcherGRPCAllowedAgentIDs != "agent-cell-a,agent-cell-b" { t.Fatalf("unified Dispatcher gRPC settings were not loaded: %+v", c) } diff --git a/internal/config/dispatcher_file.go b/internal/config/dispatcher_file.go index f823702..7dc6859 100644 --- a/internal/config/dispatcher_file.go +++ b/internal/config/dispatcher_file.go @@ -9,14 +9,12 @@ import ( "net/url" "os" "path" - "regexp" "strings" "sync" "time" "unicode/utf8" "git.ipao.vip/rogee/go-sip/contracts" - "git.ipao.vip/rogee/go-sip/internal/tenant" "github.com/santhosh-tekuri/jsonschema/v6" ) @@ -33,8 +31,6 @@ type dispatcherFile struct { } `json:"oss"` } -var credentialEnvName = regexp.MustCompile(`^[A-Za-z_][A-Za-z0-9_]*$`) - var dispatcherFileSchema = sync.OnceValues(func() (*jsonschema.Schema, error) { data, err := contracts.Files.ReadFile("upstream/" + contracts.SourceCommit + "/dispatcher-config.schema.json") if err != nil { @@ -104,12 +100,6 @@ func (c *Config) LoadDispatcherFile(filename string) error { } return errors.New("dispatcher configuration schema validation failed") } - if input.SchemaVersion != "1.0" { - return errors.New("dispatcher configuration schema_version must be 1.0") - } - if err := tenant.ValidateDispatcherID(input.DispatcherID); err != nil { - return err - } endpoint, err := url.Parse(input.OSS.Endpoint) if err != nil || endpoint.Hostname() == "" || (endpoint.Scheme != "https" && endpoint.Scheme != "http") || endpoint.User != nil || endpoint.RawQuery != "" || endpoint.Fragment != "" { return errors.New("oss.endpoint must be an HTTP(S) service URL without credentials, query or fragment") @@ -145,9 +135,6 @@ func (c *Config) LoadDispatcherFile(filename string) error { } func fileCredential(field, reference string) (string, error) { - if !credentialEnvName.MatchString(reference) { - return "", fmt.Errorf("oss.%s must name a credential environment variable", field) - } value, exists := os.LookupEnv(reference) if !exists || strings.TrimSpace(value) == "" { return "", fmt.Errorf("credential referenced by oss.%s is unavailable", field) diff --git a/internal/config/oss_file_only_test.go b/internal/config/oss_file_only_test.go index 552c74c..772065b 100644 --- a/internal/config/oss_file_only_test.go +++ b/internal/config/oss_file_only_test.go @@ -7,7 +7,10 @@ func TestLegacyOSSEnvironmentIsNotAConfigurationSource(t *testing.T) { t.Setenv(key, "legacy-value") } t.Setenv("DISPATCHER_OSS_GRANT_TTL_SECONDS", "1") - cfg := FromEnv() + cfg, err := FromEnv() + if err != nil { + t.Fatal(err) + } if cfg.OSSRegion != "" || cfg.OSSEndpoint != "" || cfg.OSSBucket != "" || cfg.OSSAccessKeyID != "" || cfg.OSSAccessKeySecret != "" || cfg.OSSKeyPrefix != "" || cfg.OSSGrantTTL != 0 { t.Fatal("legacy environment still supplies OSS settings before the required JSON file") } diff --git a/internal/dispatcher/dispatcher.go b/internal/dispatcher/dispatcher.go index b16531b..20c3dbd 100644 --- a/internal/dispatcher/dispatcher.go +++ b/internal/dispatcher/dispatcher.go @@ -155,85 +155,4 @@ func (d *Dispatcher) ApplyControl(executionID string, expectedRevision int64, ac return d.store.ApplyControl(executionID, expectedRevision, action) } -type FairScheduler struct { - mu sync.Mutex - tenants []string - cursor int -} - -func NewFairScheduler(tenants []string) *FairScheduler { - copyTenants := append([]string(nil), tenants...) - return &FairScheduler{tenants: copyTenants} -} - -func NewFairSchedulerFromStore(st *store.Store, scope string, tenants []string) (*FairScheduler, error) { - if st == nil { - return nil, errors.New("store is required") - } - scheduler := NewFairScheduler(tenants) - cursor, err := st.LoadSchedulerCursor(scope, tenants) - if err != nil { - return nil, err - } - scheduler.RestoreCursor(cursor) - return scheduler, nil -} - -// NextTenantDurable persists the next cursor before returning a tenant. A -// restart therefore resumes the bounded rotation instead of resetting to the -// first active tenant. -func (s *FairScheduler) NextTenantDurable(st *store.Store, scope string) (string, bool, error) { - if st == nil { - return "", false, errors.New("store is required") - } - s.mu.Lock() - defer s.mu.Unlock() - if len(s.tenants) == 0 { - return "", false, nil - } - index := s.cursor % len(s.tenants) - tenant := s.tenants[index] - next := (index + 1) % len(s.tenants) - if err := st.SaveSchedulerCursor(scope, s.tenants, next); err != nil { - return "", false, err - } - s.cursor = next - return tenant, true, nil -} - -// NextTenant returns the next tenant in a bounded round-robin cycle. The -// caller performs the durable task/lease check; no unbounded in-memory FIFO is -// used and a failed tenant does not consume another tenant's turn. -func (s *FairScheduler) NextTenant() (string, bool) { - s.mu.Lock() - defer s.mu.Unlock() - if len(s.tenants) == 0 { - return "", false - } - tenant := s.tenants[s.cursor%len(s.tenants)] - s.cursor = (s.cursor + 1) % len(s.tenants) - return tenant, true -} - -func (s *FairScheduler) Snapshot() (tenants []string, cursor int) { - s.mu.Lock() - defer s.mu.Unlock() - return append([]string(nil), s.tenants...), s.cursor -} - -// RestoreCursor is used only during restart recovery after the durable -// scheduler has reconstructed the active tenant set. -func (s *FairScheduler) RestoreCursor(cursor int) { - s.mu.Lock() - defer s.mu.Unlock() - if len(s.tenants) == 0 { - s.cursor = 0 - return - } - if cursor < 0 { - cursor = 0 - } - s.cursor = cursor % len(s.tenants) -} - func IsNoTask(err error) bool { return errors.Is(err, sql.ErrNoRows) } diff --git a/internal/dispatcher/dispatcher_test.go b/internal/dispatcher/dispatcher_test.go index 58ce22f..9b81235 100644 --- a/internal/dispatcher/dispatcher_test.go +++ b/internal/dispatcher/dispatcher_test.go @@ -284,55 +284,3 @@ func TestExecuteReservedBindsQuotaAndAgentExecution(t *testing.T) { t.Fatalf("task status=%q, want running", taskStatus) } } - -func TestFairSchedulerRoundRobinAndRestore(t *testing.T) { - s := NewFairScheduler([]string{"a", "b", "c"}) - for i, want := range []string{"a", "b", "c", "a"} { - got, ok := s.NextTenant() - if !ok || got != want { - t.Fatalf("turn %d = %q/%v, want %q", i, got, ok, want) - } - } - s.RestoreCursor(2) - got, _ := s.NextTenant() - if got != "c" { - t.Fatalf("restored cursor = %q, want c", got) - } -} - -func TestFairSchedulerPersistsCursorAcrossRestart(t *testing.T) { - dbPath := filepath.Join(t.TempDir(), "dispatcher.db") - st, err := store.Open(dbPath) - if err != nil { - t.Fatal(err) - } - first, err := NewFairSchedulerFromStore(st, "tenant-rotation", []string{"a", "b", "c"}) - if err != nil { - t.Fatal(err) - } - for i, want := range []string{"a", "b"} { - got, ok, err := first.NextTenantDurable(st, "tenant-rotation") - if err != nil || !ok || got != want { - t.Fatalf("turn %d = %q/%v err=%v, want %q", i, got, ok, err, want) - } - } - if err := st.Close(); err != nil { - t.Fatal(err) - } - st, err = store.Open(dbPath) - if err != nil { - t.Fatal(err) - } - defer st.Close() - if err := st.BindDispatcherID(testfixture.DispatcherID); err != nil { - t.Fatal(err) - } - second, err := NewFairSchedulerFromStore(st, "tenant-rotation", []string{"a", "b", "c"}) - if err != nil { - t.Fatal(err) - } - got, ok, err := second.NextTenantDurable(st, "tenant-rotation") - if err != nil || !ok || got != "c" { - t.Fatalf("restart turn = %q/%v err=%v, want c", got, ok, err) - } -} diff --git a/internal/store/store.go b/internal/store/store.go index 392f348..013d735 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -98,61 +98,6 @@ func (s *Store) Close() error { return s.db.Close() } func (s *Store) DB() *sql.DB { return s.db } -func (s *Store) LoadSchedulerCursor(scope string, tenants []string) (int, error) { - if scope == "" { - return 0, errors.New("scheduler scope is required") - } - encoded, err := json.Marshal(tenants) - if err != nil { - return 0, err - } - s.mu.Lock() - defer s.mu.Unlock() - var cursor int - err = s.db.QueryRow(`SELECT cursor FROM scheduler_state WHERE scope = ?`, scope).Scan(&cursor) - if errors.Is(err, sql.ErrNoRows) { - _, err = s.db.Exec(`INSERT INTO scheduler_state(scope, tenants_json, cursor, updated_at) VALUES(?, ?, 0, ?)`, scope, encoded, s.now().UTC().Format(time.RFC3339Nano)) - return 0, err - } - if err != nil { - return 0, err - } - if len(tenants) == 0 { - cursor = 0 - } else { - cursor %= len(tenants) - if cursor < 0 { - cursor += len(tenants) - } - } - _, err = s.db.Exec(`UPDATE scheduler_state SET tenants_json = ?, cursor = ?, updated_at = ? WHERE scope = ?`, encoded, cursor, s.now().UTC().Format(time.RFC3339Nano), scope) - return cursor, err -} - -func (s *Store) SaveSchedulerCursor(scope string, tenants []string, cursor int) error { - if scope == "" { - return errors.New("scheduler scope is required") - } - if len(tenants) == 0 { - cursor = 0 - } else { - cursor %= len(tenants) - if cursor < 0 { - cursor += len(tenants) - } - } - encoded, err := json.Marshal(tenants) - if err != nil { - return err - } - s.mu.Lock() - defer s.mu.Unlock() - _, err = s.db.Exec(`INSERT INTO scheduler_state(scope, tenants_json, cursor, updated_at) VALUES(?, ?, ?, ?) - ON CONFLICT(scope) DO UPDATE SET tenants_json = excluded.tenants_json, cursor = excluded.cursor, updated_at = excluded.updated_at`, - scope, encoded, cursor, s.now().UTC().Format(time.RFC3339Nano)) - return err -} - type IngestResult struct { CommandID string ExecutionID string