package outbox import ( "context" "encoding/json" "errors" "fmt" "strings" "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) var ( ErrStoreUnavailable = errors.New("outbox store unavailable") ErrEventNotFound = errors.New("outbox event not found") ) type Event struct { EventID string `json:"event_id"` EventType string `json:"event_type"` EventVersion int `json:"event_version"` TenantID *string `json:"tenant_id"` AggregateType string `json:"aggregate_type"` AggregateID string `json:"aggregate_id"` Payload json.RawMessage `json:"payload"` TraceContext json.RawMessage `json:"trace_context"` OccurredAt time.Time `json:"occurred_at"` Attempts int `json:"attempts"` } type Store struct{ pool *pgxpool.Pool } func NewStore(pool *pgxpool.Pool) *Store { return &Store{pool: pool} } func (s *Store) Claim(ctx context.Context, workerID string, limit int, lease time.Duration) ([]Event, error) { if s == nil || s.pool == nil { return nil, ErrStoreUnavailable } rows, err := s.pool.Query(ctx, ` WITH candidates AS ( SELECT event_id FROM gateway.outbox_events WHERE processed_at IS NULL AND dead_lettered_at IS NULL AND available_at <= clock_timestamp() AND (locked_at IS NULL OR locked_at < clock_timestamp()-($1 * interval '1 second')) ORDER BY available_at,occurred_at FOR UPDATE SKIP LOCKED LIMIT $2 ) UPDATE gateway.outbox_events e SET locked_at=clock_timestamp(),locked_by=$3,attempts=e.attempts+1 FROM candidates c WHERE e.event_id=c.event_id RETURNING e.event_id::text,e.event_type,e.event_version,e.tenant_id::text,e.aggregate_type, e.aggregate_id,e.payload,e.trace_context,e.occurred_at,e.attempts`, lease.Seconds(), limit, workerID) if err != nil { return nil, fmt.Errorf("%w: claim: %v", ErrStoreUnavailable, err) } defer rows.Close() events := make([]Event, 0, limit) for rows.Next() { var event Event if err := rows.Scan(&event.EventID, &event.EventType, &event.EventVersion, &event.TenantID, &event.AggregateType, &event.AggregateID, &event.Payload, &event.TraceContext, &event.OccurredAt, &event.Attempts); err != nil { return nil, fmt.Errorf("%w: scan claim: %v", ErrStoreUnavailable, err) } events = append(events, event) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("%w: claim rows: %v", ErrStoreUnavailable, err) } return events, nil } func (s *Store) MarkProcessed(ctx context.Context, eventID, workerID, streamID string) error { if s == nil || s.pool == nil { return ErrStoreUnavailable } result, err := s.pool.Exec(ctx, `UPDATE gateway.outbox_events SET processed_at=clock_timestamp(),locked_at=NULL,locked_by=NULL,last_error=NULL,published_stream_id=$3 WHERE event_id=$1 AND locked_by=$2 AND processed_at IS NULL`, eventID, workerID, streamID) if err != nil { return fmt.Errorf("%w: mark processed: %v", ErrStoreUnavailable, err) } if result.RowsAffected() != 1 { return ErrEventNotFound } return nil } func (s *Store) MarkFailed(ctx context.Context, event Event, workerID string, deliveryErr error, maxAttempts int, delay time.Duration) error { if s == nil || s.pool == nil { return ErrStoreUnavailable } message := deliveryErr.Error() if len(message) > 2048 { message = message[:2048] } dead := event.Attempts >= maxAttempts result, err := s.pool.Exec(ctx, ` UPDATE gateway.outbox_events SET locked_at=NULL,locked_by=NULL,last_error=$3, available_at=CASE WHEN $4 THEN available_at ELSE clock_timestamp()+($5 * interval '1 second') END, dead_lettered_at=CASE WHEN $4 THEN clock_timestamp() ELSE NULL END WHERE event_id=$1 AND locked_by=$2 AND processed_at IS NULL`, event.EventID, workerID, message, dead, delay.Seconds()) if err != nil { return fmt.Errorf("%w: mark failed: %v", ErrStoreUnavailable, err) } if result.RowsAffected() != 1 { return ErrEventNotFound } return nil } type EventView struct { Event AvailableAt time.Time `json:"available_at"` LockedAt *time.Time `json:"locked_at"` LockedBy *string `json:"locked_by"` ProcessedAt *time.Time `json:"processed_at"` DeadLetteredAt *time.Time `json:"dead_lettered_at"` LastError *string `json:"last_error"` PublishedStream *string `json:"published_stream_id"` } func (s *Store) List(ctx context.Context, status, eventType string, limit int) ([]EventView, error) { if s == nil || s.pool == nil { return nil, ErrStoreUnavailable } where := "TRUE" switch status { case "pending": where = "processed_at IS NULL AND dead_lettered_at IS NULL" case "dead": where = "dead_lettered_at IS NOT NULL" case "processed": where = "processed_at IS NOT NULL" } args := []any{limit} if eventType != "" { args = append(args, eventType) where += fmt.Sprintf(" AND event_type=$%d", len(args)) } rows, err := s.pool.Query(ctx, `SELECT event_id::text,event_type,event_version,tenant_id::text,aggregate_type,aggregate_id,payload,trace_context,occurred_at,attempts,available_at,locked_at,locked_by,processed_at,dead_lettered_at,last_error,published_stream_id FROM gateway.outbox_events WHERE `+where+` ORDER BY occurred_at DESC LIMIT $1`, args...) if err != nil { return nil, fmt.Errorf("%w: list: %v", ErrStoreUnavailable, err) } defer rows.Close() items := make([]EventView, 0) for rows.Next() { var item EventView if err := rows.Scan(&item.EventID, &item.EventType, &item.EventVersion, &item.TenantID, &item.AggregateType, &item.AggregateID, &item.Payload, &item.TraceContext, &item.OccurredAt, &item.Attempts, &item.AvailableAt, &item.LockedAt, &item.LockedBy, &item.ProcessedAt, &item.DeadLetteredAt, &item.LastError, &item.PublishedStream); err != nil { return nil, fmt.Errorf("%w: scan list: %v", ErrStoreUnavailable, err) } items = append(items, item) } return items, rows.Err() } func (s *Store) Retry(ctx context.Context, eventID string) error { if s == nil || s.pool == nil { return ErrStoreUnavailable } result, err := s.pool.Exec(ctx, `UPDATE gateway.outbox_events SET attempts=0,available_at=clock_timestamp(),locked_at=NULL,locked_by=NULL,dead_lettered_at=NULL,last_error=NULL WHERE event_id=$1 AND processed_at IS NULL`, eventID) if err != nil { return fmt.Errorf("%w: retry: %v", ErrStoreUnavailable, err) } if result.RowsAffected() != 1 { return ErrEventNotFound } return nil } type ConsumerHandler func(context.Context, pgx.Tx) error func (s *Store) Consume(ctx context.Context, subscriber, eventID string, handler ConsumerHandler) (bool, error) { if s == nil || s.pool == nil { return false, ErrStoreUnavailable } subscriber = strings.TrimSpace(subscriber) if subscriber == "" || eventID == "" || handler == nil { return false, errors.New("subscriber, event ID, and handler are required") } tx, err := s.pool.Begin(ctx) if err != nil { return false, fmt.Errorf("%w: begin consumption: %v", ErrStoreUnavailable, err) } defer func() { _ = tx.Rollback(ctx) }() var inserted bool err = tx.QueryRow(ctx, `WITH inserted AS (INSERT INTO gateway.event_consumptions(subscriber,event_id) VALUES($1,$2) ON CONFLICT DO NOTHING RETURNING 1) SELECT EXISTS(SELECT 1 FROM inserted)`, subscriber, eventID).Scan(&inserted) if err != nil { return false, fmt.Errorf("%w: reserve consumption: %v", ErrStoreUnavailable, err) } if !inserted { return false, nil } if err := handler(ctx, tx); err != nil { return false, err } if err := tx.Commit(ctx); err != nil { return false, fmt.Errorf("%w: commit consumption: %v", ErrStoreUnavailable, err) } return true, nil }