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 }