package store import ( "context" "database/sql" "fmt" "time" ) // OutboxEvent is the durable hand-off between a committed domain mutation and // transient WebSocket delivery. Consumers must make delivery idempotent by // event ID and only acknowledge after successful fan-out. 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 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") // 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) { if db == nil || limit < 1 || limit > 1000 { return nil, fmt.Errorf("invalid outbox read arguments") } rows, err := db.QueryContext(ctx, OutboxUnpublishedSelectSQL, 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 }