Files
2026-09-01 16:30:35 +01:00

77 lines
3.0 KiB
Go

package store
import (
"context"
"errors"
"reflect"
"testing"
"time"
)
func TestOutboxSQLPreservesReplayableOrderedReadAndPublishAck(t *testing.T) {
for query, fragments := range map[string][]string{
OutboxUnpublishedSelectSQL: {"published_at IS NULL", "ORDER BY created_at, event_id", "LIMIT $1"},
OutboxUnpublishedProposalSelectSQL: {"published_at IS NULL", "event_type = 'proposal_changed'", "ORDER BY created_at, event_id", "LIMIT $1"},
OutboxUnpublishedResultSelectSQL: {"published_at IS NULL", "event_type = 'match_completed'", "ORDER BY created_at, event_id", "LIMIT $1"},
OutboxUnpublishedStateSelectSQL: {"published_at IS NULL", "event_type = 'state_changed'", "ORDER BY created_at, event_id", "LIMIT $1"},
MatchParticipantIDsSQL: {"SELECT player_id", "match_participants", "match_id = $1", "ORDER BY player_id"},
OutboxMarkPublishedSQL: {"published_at = $2", "event_id = $1", "published_at IS NULL"},
} {
for _, fragment := range fragments {
if !contains(query, fragment) {
t.Fatalf("query %q missing %q", query, fragment)
}
}
}
}
func TestOutboxDispatcherAcknowledgesOnlyAfterDelivery(t *testing.T) {
events := []OutboxEvent{{EventID: "event-1"}, {EventID: "event-2"}}
var delivered, acknowledged []string
dispatcher := &OutboxDispatcher{
Read: func(context.Context, int) ([]OutboxEvent, error) { return events, nil },
Deliver: func(_ context.Context, event OutboxEvent) error {
delivered = append(delivered, event.EventID)
if event.EventID == "event-2" {
return errors.New("transient fan-out failure")
}
return nil
},
Ack: func(_ context.Context, eventID string, _ time.Time) error {
acknowledged = append(acknowledged, eventID)
return nil
},
}
count, err := dispatcher.Dispatch(context.Background(), 10, time.Unix(1000, 0))
if err == nil || count != 1 {
t.Fatalf("dispatch = (%d, %v), want one acknowledged event and an error", count, err)
}
if !reflect.DeepEqual(delivered, []string{"event-1", "event-2"}) || !reflect.DeepEqual(acknowledged, []string{"event-1"}) {
t.Fatalf("delivery/ack order = %v/%v", delivered, acknowledged)
}
}
func TestOutboxDispatcherRejectsInvalidConfiguration(t *testing.T) {
if count, err := (*OutboxDispatcher)(nil).Dispatch(context.Background(), 1, time.Unix(1000, 0)); err == nil || count != 0 {
t.Fatal("nil dispatcher accepted")
}
}
func TestOutboxAdaptersRejectUnsafeArgumentsWithoutDatabase(t *testing.T) {
if _, err := ReadUnpublishedOutbox(nil, nil, 1); err == nil {
t.Fatal("nil database accepted")
}
if _, err := ReadUnpublishedOutbox(nil, nil, 1001); err == nil {
t.Fatal("unbounded outbox batch accepted")
}
if err := MarkOutboxPublished(nil, nil, "event-1", time.Unix(1000, 0)); err == nil {
t.Fatal("nil database acknowledgement accepted")
}
if err := MarkOutboxPublished(nil, nil, "", time.Unix(1000, 0)); err == nil {
t.Fatal("empty event acknowledgement accepted")
}
if _, err := ReadMatchParticipantIDs(nil, nil, ""); err == nil {
t.Fatal("invalid participant read arguments accepted")
}
}