mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-10 16:04:04 +00:00
209 lines
7.2 KiB
Go
209 lines
7.2 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"fmt"
|
|
"time"
|
|
)
|
|
|
|
// OutboxEvent is the durable hand-off between a committed domain mutation and
|
|
// transient WebSocket publication. Consumers must make publication idempotent
|
|
// by event ID and only acknowledge after the local adapter accepts the event.
|
|
// Subscriber receipt is not durable; clients converge through REST recovery.
|
|
type OutboxEvent struct {
|
|
EventID string
|
|
AggregateType string
|
|
AggregateID string
|
|
Revision uint64
|
|
EventType string
|
|
Payload []byte
|
|
CreatedAt time.Time
|
|
PublishedAt *time.Time
|
|
}
|
|
|
|
const OutboxUnpublishedSelectSQL = `SELECT event_id, aggregate_type, aggregate_id, revision,
|
|
event_type, payload, created_at, published_at
|
|
FROM outbox
|
|
WHERE published_at IS NULL
|
|
ORDER BY created_at, event_id
|
|
LIMIT $1`
|
|
|
|
const OutboxUnpublishedProposalSelectSQL = `SELECT event_id, aggregate_type, aggregate_id, revision,
|
|
event_type, payload, created_at, published_at
|
|
FROM outbox
|
|
WHERE published_at IS NULL AND event_type = 'proposal_changed'
|
|
ORDER BY created_at, event_id
|
|
LIMIT $1`
|
|
|
|
const OutboxUnpublishedResultSelectSQL = `SELECT event_id, aggregate_type, aggregate_id, revision,
|
|
event_type, payload, created_at, published_at
|
|
FROM outbox
|
|
WHERE published_at IS NULL AND event_type = 'match_completed'
|
|
ORDER BY created_at, event_id
|
|
LIMIT $1`
|
|
|
|
const OutboxUnpublishedStateSelectSQL = `SELECT event_id, aggregate_type, aggregate_id, revision,
|
|
event_type, payload, created_at, published_at
|
|
FROM outbox
|
|
WHERE published_at IS NULL AND event_type = 'state_changed'
|
|
ORDER BY created_at, event_id
|
|
LIMIT $1`
|
|
|
|
const MatchParticipantIDsSQL = `SELECT player_id
|
|
FROM match_participants
|
|
WHERE match_id = $1
|
|
ORDER BY player_id`
|
|
|
|
const OutboxMarkPublishedSQL = `UPDATE outbox
|
|
SET published_at = $2
|
|
WHERE event_id = $1 AND published_at IS NULL`
|
|
|
|
var ErrOutboxEventNotFound = fmt.Errorf("outbox event not found or already published")
|
|
|
|
type OutboxDelivery func(context.Context, OutboxEvent) error
|
|
|
|
// OutboxDispatcher is the durable-to-transient bridge. Read and Ack are
|
|
// injectable so ordering can be tested without a live PostgreSQL instance.
|
|
// Adapter invocation is at-least-once: a crash after invocation and before
|
|
// acknowledgement leaves the event replayable, while an adapter failure stops
|
|
// the batch. This does not imply that a transient subscriber received it.
|
|
type OutboxDispatcher struct {
|
|
Read func(context.Context, int) ([]OutboxEvent, error)
|
|
Ack func(context.Context, string, time.Time) error
|
|
Deliver OutboxDelivery
|
|
}
|
|
|
|
func NewOutboxDispatcher(db *sql.DB, deliver OutboxDelivery) *OutboxDispatcher {
|
|
return &OutboxDispatcher{
|
|
Read: func(ctx context.Context, limit int) ([]OutboxEvent, error) {
|
|
return ReadUnpublishedOutbox(ctx, db, limit)
|
|
},
|
|
Ack: func(ctx context.Context, eventID string, publishedAt time.Time) error {
|
|
return MarkOutboxPublished(ctx, db, eventID, publishedAt)
|
|
},
|
|
Deliver: deliver,
|
|
}
|
|
}
|
|
|
|
// Dispatch publishes at most limit events in the store's stable order and
|
|
// returns the number acknowledged. A successful delivery followed by an ack
|
|
// error intentionally leaves that event replayable.
|
|
func (d *OutboxDispatcher) Dispatch(ctx context.Context, limit int, publishedAt time.Time) (int, error) {
|
|
if d == nil || d.Read == nil || d.Ack == nil || d.Deliver == nil || limit < 1 || limit > 1000 || publishedAt.IsZero() {
|
|
return 0, fmt.Errorf("invalid outbox dispatcher")
|
|
}
|
|
events, err := d.Read(ctx, limit)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
acknowledged := 0
|
|
for _, event := range events {
|
|
if event.EventID == "" {
|
|
return acknowledged, fmt.Errorf("outbox event has no ID")
|
|
}
|
|
if err := d.Deliver(ctx, event); err != nil {
|
|
return acknowledged, err
|
|
}
|
|
if err := d.Ack(ctx, event.EventID, publishedAt); err != nil {
|
|
return acknowledged, err
|
|
}
|
|
acknowledged++
|
|
}
|
|
return acknowledged, nil
|
|
}
|
|
|
|
// ReadUnpublishedOutbox returns a bounded, stable ordered batch. It does not
|
|
// mark rows before delivery: a worker crash therefore leaves events replayable.
|
|
func ReadUnpublishedOutbox(ctx context.Context, db *sql.DB, limit int) ([]OutboxEvent, error) {
|
|
return readUnpublishedOutbox(ctx, db, limit, OutboxUnpublishedSelectSQL)
|
|
}
|
|
|
|
// ReadUnpublishedProposalOutbox returns only WebSocket-routable proposal
|
|
// events. Other outbox consumers (for example result reconciliation) retain
|
|
// ownership of their event types and cannot be acknowledged accidentally by
|
|
// the control-plane WebSocket dispatcher.
|
|
func ReadUnpublishedProposalOutbox(ctx context.Context, db *sql.DB, limit int) ([]OutboxEvent, error) {
|
|
return readUnpublishedOutbox(ctx, db, limit, OutboxUnpublishedProposalSelectSQL)
|
|
}
|
|
|
|
// ReadUnpublishedResultOutbox returns only durable match-completion events.
|
|
// Proposal and result consumers acknowledge separate event types so one
|
|
// transient fan-out outage cannot hide rows owned by another consumer.
|
|
func ReadUnpublishedResultOutbox(ctx context.Context, db *sql.DB, limit int) ([]OutboxEvent, error) {
|
|
return readUnpublishedOutbox(ctx, db, limit, OutboxUnpublishedResultSelectSQL)
|
|
}
|
|
|
|
// ReadUnpublishedStateOutbox returns lifecycle state events, leaving proposal
|
|
// and result rows to their dedicated consumers.
|
|
func ReadUnpublishedStateOutbox(ctx context.Context, db *sql.DB, limit int) ([]OutboxEvent, error) {
|
|
return readUnpublishedOutbox(ctx, db, limit, OutboxUnpublishedStateSelectSQL)
|
|
}
|
|
|
|
func ReadMatchParticipantIDs(ctx context.Context, db *sql.DB, matchID string) ([]string, error) {
|
|
if db == nil || matchID == "" {
|
|
return nil, fmt.Errorf("invalid match participant read arguments")
|
|
}
|
|
rows, err := db.QueryContext(ctx, MatchParticipantIDsSQL, matchID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var players []string
|
|
for rows.Next() {
|
|
var playerID string
|
|
if err := rows.Scan(&playerID); err != nil {
|
|
return nil, err
|
|
}
|
|
players = append(players, playerID)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return players, nil
|
|
}
|
|
|
|
func readUnpublishedOutbox(ctx context.Context, db *sql.DB, limit int, query string) ([]OutboxEvent, error) {
|
|
if db == nil || limit < 1 || limit > 1000 {
|
|
return nil, fmt.Errorf("invalid outbox read arguments")
|
|
}
|
|
rows, err := db.QueryContext(ctx, query, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
events := make([]OutboxEvent, 0, limit)
|
|
for rows.Next() {
|
|
var event OutboxEvent
|
|
if err := rows.Scan(&event.EventID, &event.AggregateType, &event.AggregateID, &event.Revision, &event.EventType, &event.Payload, &event.CreatedAt, &event.PublishedAt); err != nil {
|
|
return nil, err
|
|
}
|
|
events = append(events, event)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return events, nil
|
|
}
|
|
|
|
// MarkOutboxPublished acknowledges one event only if it is still unpublished.
|
|
// Repeated acknowledgement is reported to the caller so a worker cannot
|
|
// mistake an already-completed delivery for a fresh one.
|
|
func MarkOutboxPublished(ctx context.Context, db *sql.DB, eventID string, publishedAt time.Time) error {
|
|
if db == nil || eventID == "" || publishedAt.IsZero() {
|
|
return fmt.Errorf("invalid outbox acknowledgement arguments")
|
|
}
|
|
result, err := db.ExecContext(ctx, OutboxMarkPublishedSQL, eventID, publishedAt)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
changed, err := result.RowsAffected()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if changed != 1 {
|
|
return ErrOutboxEventNotFound
|
|
}
|
|
return nil
|
|
}
|