124 lines
3.5 KiB
Go
124 lines
3.5 KiB
Go
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)
|
|
}
|