84 lines
2.4 KiB
Go
84 lines
2.4 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 TestMQControlRejectsLateRevisionWithoutRegressingTask(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)
|
|
}
|
|
first, err := contracts.Files.ReadFile(base + "task-control.json")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, duplicate, err := s.HandleTaskControl(first, route.InboundKey); err != nil || duplicate {
|
|
t.Fatalf("first control: err=%v duplicate=%v", err, duplicate)
|
|
}
|
|
|
|
var late map[string]any
|
|
if err := json.Unmarshal(first, &late); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
late["command_id"] = "task.control-late"
|
|
payload := late["payload"].(map[string]any)
|
|
payload["action"] = "resume"
|
|
payload["expected_task_revision"] = float64(1)
|
|
lateRaw, err := json.Marshal(late)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
responseID, duplicate, err := s.HandleTaskControl(lateRaw, route.InboundKey)
|
|
if err != nil || duplicate {
|
|
t.Fatalf("late control: err=%v duplicate=%v", err, duplicate)
|
|
}
|
|
var body []byte
|
|
if err := s.DB().QueryRow(`SELECT body FROM outbox WHERE event_id=?`, responseID).Scan(&body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var event struct {
|
|
Payload struct {
|
|
Status string `json:"status"`
|
|
Reason string `json:"reason_code"`
|
|
} `json:"payload"`
|
|
}
|
|
if err := json.Unmarshal(body, &event); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if event.Payload.Status != "rejected" || event.Payload.Reason != "revision_conflict" {
|
|
t.Fatalf("late control was not rejected: %+v", event.Payload)
|
|
}
|
|
var status string
|
|
var revision int
|
|
if err := s.DB().QueryRow(`SELECT status,task_revision FROM tasks`).Scan(&status, &revision); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if status != "paused" || revision != 2 {
|
|
t.Fatalf("late control regressed task: status=%s revision=%d", status, revision)
|
|
}
|
|
}
|