169 lines
5.4 KiB
Go
169 lines
5.4 KiB
Go
package store
|
|
|
|
import (
|
|
"database/sql"
|
|
"encoding/json"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"testing"
|
|
|
|
"git.ipao.vip/rogee/go-sip/internal/configread"
|
|
)
|
|
|
|
const currentDispatcherID = "c046b893-8628-4589-ae50-619d049248a6"
|
|
|
|
func currentStoreSnapshot(t *testing.T) configread.Snapshot {
|
|
t.Helper()
|
|
dir := filepath.Join("..", "..", "contracts", "local", "examples")
|
|
read := func(name string, dst any) {
|
|
t.Helper()
|
|
raw, err := os.ReadFile(filepath.Join(dir, name+".json"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := json.Unmarshal(raw, dst); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
var s configread.Snapshot
|
|
read("config-read-task-asr", &s.Task)
|
|
read("config-read-sip", &s.SIP)
|
|
read("config-read-quota", &s.Quota)
|
|
var p struct {
|
|
Providers []configread.Provider `json:"providers"`
|
|
}
|
|
read("config-read-providers", &p)
|
|
s.Providers = make(map[string]configread.Provider)
|
|
for _, provider := range p.Providers {
|
|
s.Providers[provider.ProviderRef] = provider
|
|
}
|
|
return s
|
|
}
|
|
|
|
func TestStoreFreshSnapshotAndControlSurviveRestart(t *testing.T) {
|
|
path := filepath.Join(t.TempDir(), "current.db")
|
|
s, err := Open(path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
task := configread.DiscoveredTask{TaskID: "task-asr", TenantID: 1001, TaskRevision: 1, Status: "running"}
|
|
if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{task}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if allowed, _ := s.CanAdmit(currentDispatcherID, 1001, "task-asr"); allowed {
|
|
t.Fatal("opened admission before control backlog drained")
|
|
}
|
|
if err := s.SaveSnapshot(currentStoreSnapshot(t)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.MarkReadyForSIP(currentDispatcherID, 8); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if allowed, err := s.CanAdmit(currentDispatcherID, 1001, "task-asr"); err != nil || !allowed {
|
|
t.Fatalf("fresh running task not admitted: %v, %v", allowed, err)
|
|
}
|
|
if err := s.ApplyControl(currentDispatcherID, 1001, "task-asr", "pause"); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
s, err = Open(path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer s.Close()
|
|
if allowed, _ := s.CanAdmit(currentDispatcherID, 1001, "task-asr"); allowed {
|
|
t.Fatal("restart erased durable pause")
|
|
}
|
|
if err := s.ApplyDiscoverySnapshot(currentDispatcherID, []configread.DiscoveredTask{task}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.MarkReadyForSIP(currentDispatcherID, 8); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if allowed, _ := s.CanAdmit(currentDispatcherID, 1001, "task-asr"); allowed {
|
|
t.Fatal("running discovery undid durable pause")
|
|
}
|
|
if err := s.ApplyControl(currentDispatcherID, 1001, "task-asr", "stop"); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.ApplyControl(currentDispatcherID, 1001, "task-asr", "resume"); err == nil {
|
|
t.Fatal("stopped task resumed")
|
|
}
|
|
}
|
|
|
|
func TestStoreRejectsSameRevisionDifferentContent(t *testing.T) {
|
|
s, err := Open(filepath.Join(t.TempDir(), "current.db"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer s.Close()
|
|
snapshot := currentStoreSnapshot(t)
|
|
if err := s.SaveSnapshot(snapshot); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
snapshot.Task.Name = "changed without changing revision"
|
|
snapshot.Task.Raw = []byte(strings.Replace(string(snapshot.Task.Raw), `"ASR example"`, `"changed without changing revision"`, 1))
|
|
if err := s.SaveSnapshot(snapshot); err == nil {
|
|
t.Fatal("accepted different content under same immutable revision")
|
|
}
|
|
}
|
|
|
|
func TestStoreRefusesOldSchemaWithoutDeletingRows(t *testing.T) {
|
|
path := filepath.Join(t.TempDir(), "legacy.db")
|
|
legacy, err := sql.Open("sqlite", path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := legacy.Exec(`CREATE TABLE local_v01_outbox(id INTEGER PRIMARY KEY, tenant_key TEXT NOT NULL); INSERT INTO local_v01_outbox(tenant_key) VALUES('not-a-number')`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
legacy.Close()
|
|
if s, err := Open(path); err == nil {
|
|
s.Close()
|
|
t.Fatal("opened an old database without an explicit data decision")
|
|
}
|
|
legacy, err = sql.Open("sqlite", path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer legacy.Close()
|
|
var tenantKey string
|
|
if err := legacy.QueryRow(`SELECT tenant_key FROM local_v01_outbox`).Scan(&tenantKey); err != nil || tenantKey != "not-a-number" {
|
|
t.Fatalf("old data changed or disappeared: %q, %v", tenantKey, err)
|
|
}
|
|
}
|
|
|
|
func TestStoreRefusesPreviousCurrentLayoutBeforeModifyingDatabase(t *testing.T) {
|
|
path := filepath.Join(t.TempDir(), "previous-current.db")
|
|
old, err := sql.Open("sqlite", path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := old.Exec(`CREATE TABLE dispatcher_state(dispatcher_id TEXT PRIMARY KEY,discovery_ready INTEGER NOT NULL); INSERT INTO dispatcher_state VALUES('original-dispatcher',1)`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := old.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if opened, err := Open(path); err == nil {
|
|
opened.Close()
|
|
t.Fatal("accepted prior current SQLite layout without explicit data decision")
|
|
}
|
|
old, err = sql.Open("sqlite", path)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer old.Close()
|
|
var id string
|
|
var ready, count int
|
|
if err := old.QueryRow(`SELECT dispatcher_id,discovery_ready FROM dispatcher_state`).Scan(&id, &ready); err != nil || id != "original-dispatcher" || ready != 1 {
|
|
t.Fatalf("old admission state changed: %q %d %v", id, ready, err)
|
|
}
|
|
if err := old.QueryRow(`SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%'`).Scan(&count); err != nil || count != 1 {
|
|
t.Fatalf("old database was modified before rejection: tables=%d err=%v", count, err)
|
|
}
|
|
}
|