package outbox import ( "context" "time" "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/v2/bson" "go.mongodb.org/mongo-driver/v2/mongo" "go.uber.org/zap" ) type mongoStore struct { logger mlogger.Logger repo repository.Repository } func NewMongoStore(logger mlogger.Logger, db *mongo.Database) (Store, error) { if db == nil { return nil, merrors.InvalidArgument("mongo database is nil") } if logger == nil { logger = zap.NewNop() } repo := repository.CreateMongoRepository(db, Collection) 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 } 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(Collection) childLogger.Debug("Outbox store initialised", zap.String("collection", Collection)) return &mongoStore{logger: childLogger, repo: repo}, nil } func (o *mongoStore) Create(ctx context.Context, event *Event) error { if event == nil { o.logger.Warn("Attempt to create nil outbox event") return merrors.InvalidArgument("outbox: nil event") } if err := o.repo.Insert(ctx, event, nil); err != nil { if mongo.IsDuplicateKeyError(err) { o.logger.Warn("Duplicate outbox event id", zap.String("event_id", 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("event_id", event.EventID), zap.String("subject", event.Subject)) return nil } func (o *mongoStore) ListPending(ctx context.Context, limit int) ([]*Event, error) { limit64 := int64(limit) query := repository.Query(). Filter(repository.Field("status"), StatusPending). Limit(&limit64). Sort(repository.Field("createdAt"), true) events := make([]*Event, 0) err := o.repo.FindManyByFilter(ctx, query, func(cur *mongo.Cursor) error { doc := &Event{} 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 } return events, nil } func (o *mongoStore) MarkSent(ctx context.Context, eventRef bson.ObjectID, sentAt time.Time) error { if eventRef.IsZero() { return merrors.InvalidArgument("outbox: zero event id") } patch := repository.Patch(). Set(repository.Field("status"), StatusSent). Set(repository.Field("sentAt"), sentAt) return o.repo.Patch(ctx, eventRef, patch) } func (o *mongoStore) MarkFailed(ctx context.Context, eventRef bson.ObjectID) error { if eventRef.IsZero() { return merrors.InvalidArgument("outbox: zero event id") } patch := repository.Patch().Set(repository.Field("status"), StatusFailed) return o.repo.Patch(ctx, eventRef, patch) } func (o *mongoStore) IncrementAttempts(ctx context.Context, eventRef bson.ObjectID) error { if eventRef.IsZero() { return merrors.InvalidArgument("outbox: zero event id") } patch := repository.Patch().Inc(repository.Field("attempts"), 1) return o.repo.Patch(ctx, eventRef, patch) }