249 lines
7.9 KiB
Go
249 lines
7.9 KiB
Go
package store
|
|
|
|
import (
|
|
"encoding/json"
|
|
"errors"
|
|
"path/filepath"
|
|
"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 TestCommandQueryUsesRecordedAggregateVersion(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) }
|
|
raw, err := contracts.Files.ReadFile("upstream/v1/examples/command-query.json")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
request, err := contract.DecodeService(raw)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var target struct {
|
|
CommandID string `json:"command_id"`
|
|
}
|
|
if err := json.Unmarshal(request.Payload, &target); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.DB().Exec(`INSERT INTO inbox(tenant_id,command_id,tenant_key,command_type,body_hash,body,status,received_at,persisted_at) VALUES(?,?,?,'call.execute','fixture','{}','persisted',?,?)`, request.TenantID, target.CommandID, request.TenantKey, request.IssuedAt, request.IssuedAt); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
eventRaw, err := contracts.Files.ReadFile("upstream/v1/examples/event-command-result.json")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var event map[string]any
|
|
if err := json.Unmarshal(eventRaw, &event); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
event["aggregate_id"], event["aggregate_version"] = target.CommandID, 7
|
|
event["payload"].(map[string]any)["command_id"] = target.CommandID
|
|
eventRaw, err = json.Marshal(event)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
route, err := tenant.NewDispatcherRoute(identityA, request.TenantKey)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.DB().Exec(`INSERT INTO outbox(event_id,tenant_key,exchange,routing_key,body,status,created_at) VALUES('fact',?,'agent-call.saas.v2',?,?,'published',?)`, request.TenantKey, route.OutboundKey, eventRaw, request.IssuedAt); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
id, _, err := s.HandleQuery(raw, route.InboundKey)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var responseRaw []byte
|
|
if err := s.DB().QueryRow(`SELECT body FROM outbox WHERE event_id=?`, id).Scan(&responseRaw); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
response, err := contract.DecodeService(responseRaw)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var snapshot struct {
|
|
AggregateVersion int64 `json:"aggregate_version"`
|
|
CommandID string `json:"command_id"`
|
|
Status string `json:"status"`
|
|
}
|
|
if err := json.Unmarshal(response.Payload, &snapshot); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if response.Status != "ok" || snapshot.AggregateVersion != 7 || snapshot.CommandID != target.CommandID || snapshot.Status != "accepted" {
|
|
t.Fatalf("snapshot was fabricated or bound incorrectly: %+v", snapshot)
|
|
}
|
|
if _, err := s.DB().Exec(`DELETE FROM outbox WHERE event_id='fact'`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// Missing aggregate evidence must not be replaced by a hard-coded version.
|
|
var next map[string]any
|
|
if err := json.Unmarshal(raw, &next); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
next["message_id"] = "query-missing-evidence"
|
|
raw, err = json.Marshal(next)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, _, err := s.HandleQuery(raw, route.InboundKey); err == nil {
|
|
t.Fatal("missing aggregate evidence silently replaced")
|
|
}
|
|
}
|
|
|
|
func TestCommandQueryRecoveryAndAtomicFailure(t *testing.T) {
|
|
filename := filepath.Join(t.TempDir(), "queries.db")
|
|
s, err := Open(filename)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
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) }
|
|
raw, err := contracts.Files.ReadFile("upstream/v1/examples/command-query.json")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
route, err := tenant.NewDispatcherRoute(identityA, "tenant-a")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := s.DB().Exec(`CREATE TRIGGER reject_query_response BEFORE INSERT ON outbox BEGIN SELECT RAISE(ABORT,'injected storage failure'); END`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, _, err := s.HandleQuery(raw, route.InboundKey); err == nil {
|
|
t.Fatal("storage failure hidden")
|
|
}
|
|
for _, table := range []string{"outbox", "mq_query_inbox", "tenant_bindings"} {
|
|
var n int
|
|
if err := s.DB().QueryRow("SELECT COUNT(*) FROM " + table).Scan(&n); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if n != 0 {
|
|
t.Fatalf("partial commit in %s", table)
|
|
}
|
|
}
|
|
if _, err := s.DB().Exec(`DROP TRIGGER reject_query_response`); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
id, _, err := s.HandleQuery(raw, route.InboundKey)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := s.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
s, err = Open(filename)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer s.Close()
|
|
// A later duplicate returns the original decision, even after its deadline.
|
|
s.now = func() time.Time { return time.Date(2026, 9, 22, 0, 0, 0, 0, time.UTC) }
|
|
got, duplicate, err := s.HandleQuery(raw, route.InboundKey)
|
|
if err != nil || !duplicate || got != id {
|
|
t.Fatalf("restart lost response identity: %s %v %v", got, duplicate, err)
|
|
}
|
|
var changed map[string]any
|
|
if err := json.Unmarshal(raw, &changed); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
changed["message_id"] = "expired-query"
|
|
raw, err = json.Marshal(changed)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
expired, _, err := s.HandleQuery(raw, route.InboundKey)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var body []byte
|
|
if err := s.DB().QueryRow(`SELECT body FROM outbox WHERE event_id=?`, expired).Scan(&body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
response, err := contract.DecodeService(body)
|
|
if err != nil || response.ReasonCode != "expired" {
|
|
t.Fatalf("late query not rejected: %v %v", response, err)
|
|
}
|
|
}
|
|
|
|
func TestCommandQueryPersistsOneCorrelatedResponse(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) }
|
|
raw, err := contracts.Files.ReadFile("upstream/v1/examples/command-query.json")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
route, err := tenant.NewDispatcherRoute(identityA, "tenant-a")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
id, duplicate, err := s.HandleQuery(raw, route.InboundKey)
|
|
if err != nil || duplicate || id == "" {
|
|
t.Fatalf("first query: %s %v %v", id, duplicate, err)
|
|
}
|
|
if _, err := s.DB().Exec(`UPDATE outbox SET status='published' WHERE event_id=?`, id); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
second, duplicate, err := s.HandleQuery(raw, route.InboundKey)
|
|
if err != nil || !duplicate || second != id {
|
|
t.Fatalf("duplicate query: %s %v %v", second, duplicate, err)
|
|
}
|
|
var delivery string
|
|
if err := s.DB().QueryRow(`SELECT status FROM outbox WHERE event_id=?`, id).Scan(&delivery); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if delivery != "pending" {
|
|
t.Fatalf("duplicate cannot recover original reply: %s", delivery)
|
|
}
|
|
var body []byte
|
|
var count int
|
|
if err := s.DB().QueryRow(`SELECT COUNT(*) FROM outbox`).Scan(&count); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if count != 1 {
|
|
t.Fatalf("duplicate produced %d responses", count)
|
|
}
|
|
if err := s.DB().QueryRow(`SELECT body FROM outbox WHERE event_id=?`, id).Scan(&body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
response, err := contract.DecodeService(body)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if response.MessageType != "command.query.result" || response.CorrelationID != "command.query-a" || response.ReasonCode != "not_found" || response.DispatcherID != identityA {
|
|
t.Fatalf("wrong response identity/status: %+v", response)
|
|
}
|
|
var changed map[string]any
|
|
if err := json.Unmarshal(raw, &changed); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
changed["payload"] = map[string]any{"command_id": "different"}
|
|
changedRaw, _ := json.Marshal(changed)
|
|
if _, _, err := s.HandleQuery(changedRaw, route.InboundKey); !errors.Is(err, ErrIdempotencyConflict) {
|
|
t.Fatalf("different content accepted: %v", err)
|
|
}
|
|
other, _ := tenant.NewDispatcherRoute(identityB, "tenant-a")
|
|
if _, _, err := s.HandleQuery(raw, other.InboundKey); err == nil {
|
|
t.Fatal("foreign route accepted")
|
|
}
|
|
}
|