refactor: remove unused scheduling and cache code
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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"})
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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 }
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
+11
-22
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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) }
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user