package store import ( "context" "time" "github.com/tech/sendico/ledger/storage" "github.com/tech/sendico/ledger/storage/model" "github.com/tech/sendico/pkg/db/repository" ri "github.com/tech/sendico/pkg/db/repository/index" "github.com/tech/sendico/pkg/merrors" "github.com/tech/sendico/pkg/mlogger" "go.mongodb.org/mongo-driver/bson/primitive" "go.mongodb.org/mongo-driver/mongo" "go.uber.org/zap" ) type outboxStore struct { logger mlogger.Logger repo repository.Repository } func NewOutbox(logger mlogger.Logger, db *mongo.Database) (storage.OutboxStore, error) { repo := repository.CreateMongoRepository(db, model.OutboxCollection) // Create index on status + createdAt for efficient pending query statusIndex := &ri.Definition{ Keys: []ri.Key{ {Field: "status", Sort: ri.Asc}, {Field: "createdAt", Sort: ri.Asc}, }, } if err := repo.CreateIndex(statusIndex); err != nil { logger.Error("failed to ensure outbox status index", zap.Error(err)) return nil, err } // Create unique index on eventId for deduplication eventIdIndex := &ri.Definition{ Keys: []ri.Key{ {Field: "eventId", Sort: ri.Asc}, }, Unique: true, } if err := repo.CreateIndex(eventIdIndex); err != nil { logger.Error("failed to ensure outbox eventId index", zap.Error(err)) return nil, err } childLogger := logger.Named(model.OutboxCollection) childLogger.Debug("outbox store initialised", zap.String("collection", model.OutboxCollection)) return &outboxStore{ logger: childLogger, repo: repo, }, nil } func (o *outboxStore) Create(ctx context.Context, event *model.OutboxEvent) error { if event == nil { o.logger.Warn("attempt to create nil outbox event") return merrors.InvalidArgument("outboxStore: nil outbox event") } if err := o.repo.Insert(ctx, event, nil); err != nil { if mongo.IsDuplicateKeyError(err) { o.logger.Warn("duplicate event ID", zap.String("eventId", event.EventID)) return merrors.DataConflict("outbox event with this ID already exists") } o.logger.Warn("failed to create outbox event", zap.Error(err)) return err } o.logger.Debug("outbox event created", zap.String("eventId", event.EventID), zap.String("subject", event.Subject)) return nil } func (o *outboxStore) ListPending(ctx context.Context, limit int) ([]*model.OutboxEvent, error) { limit64 := int64(limit) query := repository.Query(). Filter(repository.Field("status"), model.OutboxStatusPending). Limit(&limit64). Sort(repository.Field("createdAt"), true) // true = ascending (oldest first) events := make([]*model.OutboxEvent, 0) err := o.repo.FindManyByFilter(ctx, query, func(cur *mongo.Cursor) error { doc := &model.OutboxEvent{} if err := cur.Decode(doc); err != nil { return err } events = append(events, doc) return nil }) if err != nil { o.logger.Warn("failed to list pending outbox events", zap.Error(err)) return nil, err } o.logger.Debug("listed pending outbox events", zap.Int("count", len(events))) return events, nil } func (o *outboxStore) MarkSent(ctx context.Context, eventRef primitive.ObjectID, sentAt time.Time) error { if eventRef.IsZero() { o.logger.Warn("attempt to mark sent with zero event ID") return merrors.InvalidArgument("outboxStore: zero event ID") } patch := repository.Patch(). Set(repository.Field("status"), model.OutboxStatusSent). Set(repository.Field("sentAt"), sentAt) if err := o.repo.Patch(ctx, eventRef, patch); err != nil { o.logger.Warn("failed to mark outbox event as sent", zap.Error(err), zap.String("eventRef", eventRef.Hex())) return err } o.logger.Debug("outbox event marked as sent", zap.String("eventRef", eventRef.Hex())) return nil } func (o *outboxStore) MarkFailed(ctx context.Context, eventRef primitive.ObjectID) error { if eventRef.IsZero() { o.logger.Warn("attempt to mark failed with zero event ID") return merrors.InvalidArgument("outboxStore: zero event ID") } patch := repository.Patch().Set(repository.Field("status"), model.OutboxStatusFailed) if err := o.repo.Patch(ctx, eventRef, patch); err != nil { o.logger.Warn("failed to mark outbox event as failed", zap.Error(err), zap.String("eventRef", eventRef.Hex())) return err } o.logger.Debug("outbox event marked as failed", zap.String("eventRef", eventRef.Hex())) return nil } func (o *outboxStore) IncrementAttempts(ctx context.Context, eventRef primitive.ObjectID) error { if eventRef.IsZero() { o.logger.Warn("attempt to increment attempts with zero event ID") return merrors.InvalidArgument("outboxStore: zero event ID") } patch := repository.Patch().Inc(repository.Field("attempts"), 1) if err := o.repo.Patch(ctx, eventRef, patch); err != nil { o.logger.Warn("failed to increment outbox attempts", zap.Error(err), zap.String("eventRef", eventRef.Hex())) return err } o.logger.Debug("outbox attempts incremented", zap.String("eventRef", eventRef.Hex())) return nil }