Files

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")
}
}