115 lines
3.5 KiB
Go
115 lines
3.5 KiB
Go
package store
|
|
|
|
import (
|
|
"encoding/json"
|
|
"testing"
|
|
"time"
|
|
|
|
"git.ipao.vip/rogee/go-sip/contracts"
|
|
"git.ipao.vip/rogee/go-sip/internal/contract"
|
|
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
|
)
|
|
|
|
func TestMQControlAppliesQueuedTasksButKeepsActiveTargetsPending(t *testing.T) {
|
|
for _, initialState := range []string{"accepted", "reserved", "running"} {
|
|
active := initialState == "running"
|
|
t.Run(initialState, func(t *testing.T) {
|
|
s, err := Open(":memory:")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer s.Close()
|
|
if err := s.BindDispatcherID(identityA); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
s.now = func() time.Time { return time.Date(2026, 9, 21, 0, 0, 1, 0, time.UTC) }
|
|
route, err := tenant.NewDispatcherRoute(identityA, "tenant-a")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
base := "upstream/" + contract.MQSourceCommit + "/examples/"
|
|
execute, err := contracts.Files.ReadFile(base + "call-execute.json")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.IngestCommand(execute, route.InboundKey); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if initialState == "reserved" {
|
|
_, payload, err := contract.DecodeExecute(execute)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.SetQuota("global", 1); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Reserve("not-yet-assigned", payload.ExecutionID, "tenant-a", []string{"global"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
} else if initialState == "running" {
|
|
if _, err := s.DB().Exec(`UPDATE tasks SET status='running'`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
raw, err := contracts.Files.ReadFile(base + "task-control.json")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
receipt, duplicate, err := s.HandleTaskControl(raw, route.InboundKey)
|
|
if err != nil || duplicate {
|
|
t.Fatalf("control: %v duplicate=%v", err, duplicate)
|
|
}
|
|
again, duplicate, err := s.HandleTaskControl(raw, route.InboundKey)
|
|
if err != nil || !duplicate || again != receipt {
|
|
t.Fatal("duplicate control did not preserve original receipt")
|
|
}
|
|
var body []byte
|
|
if err := s.DB().QueryRow(`SELECT body FROM outbox WHERE event_id=?`, receipt).Scan(&body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := contract.ValidateMQMessage(body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var event struct {
|
|
Payload struct {
|
|
Status string `json:"status"`
|
|
} `json:"payload"`
|
|
}
|
|
if err := json.Unmarshal(body, &event); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
expected := "applied"
|
|
if active {
|
|
expected = "accepted"
|
|
}
|
|
if event.Payload.Status != expected {
|
|
t.Fatalf("status=%s want=%s", event.Payload.Status, expected)
|
|
}
|
|
var taskState string
|
|
var revision int
|
|
if err := s.DB().QueryRow(`SELECT status,task_revision FROM tasks`).Scan(&taskState, &revision); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if taskState != "paused" {
|
|
t.Fatal("pause failed to close local admission")
|
|
}
|
|
if initialState == "reserved" {
|
|
var reservationState string
|
|
var reserved int
|
|
if err := s.DB().QueryRow(`SELECT state FROM reservations WHERE reservation_id='not-yet-assigned'`).Scan(&reservationState); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.DB().QueryRow(`SELECT reserved_value FROM quotas WHERE scope='global'`).Scan(&reserved); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if reservationState != "released" || reserved != 0 {
|
|
t.Fatal("never-dispatched reservation leaked after control")
|
|
}
|
|
}
|
|
if (!active && revision != 2) || (active && revision != 1) {
|
|
t.Fatal("revision pretended remote application")
|
|
}
|
|
})
|
|
}
|
|
}
|