125 lines
4.4 KiB
Go
125 lines
4.4 KiB
Go
package outbox
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
"unicode/utf8"
|
|
|
|
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
)
|
|
|
|
var (
|
|
ErrInvalidInput = errors.New("可靠投递请求不符合要求")
|
|
ErrNotFound = errors.New("可靠投递记录不存在")
|
|
ErrNotDead = errors.New("仅死信记录可以重新排队")
|
|
ErrVersionConflict = errors.New("记录已变化,请刷新后重试")
|
|
)
|
|
|
|
type Service struct {
|
|
coreService.Service
|
|
Now func() time.Time
|
|
}
|
|
|
|
func NewService(db *gorm.DB) *Service {
|
|
return &Service{Service: coreService.Service{Orm: db}, Now: time.Now}
|
|
}
|
|
|
|
func (s *Service) List(request PageRequest) (PageResponse, error) {
|
|
request.State = strings.TrimSpace(request.State)
|
|
request.InternalType = strings.TrimSpace(request.InternalType)
|
|
request.Keyword = strings.TrimSpace(request.Keyword)
|
|
if request.GetPageSize() > 100 || utf8.RuneCountInString(request.Keyword) > 128 || len(request.InternalType) > 64 || (request.State != "" && !validState(request.State)) {
|
|
return PageResponse{}, ErrInvalidInput
|
|
}
|
|
query := s.Orm.Model(&Message{})
|
|
if request.State != "" {
|
|
query = query.Where("state = ?", request.State)
|
|
}
|
|
if request.InternalType != "" {
|
|
query = query.Where("internal_type = ?", request.InternalType)
|
|
}
|
|
if request.Keyword != "" {
|
|
pattern := "%" + strings.ToLower(request.Keyword) + "%"
|
|
query = query.Where("LOWER(id) LIKE ? OR LOWER(idempotency_key) LIKE ? OR LOWER(business_ref) LIKE ?", pattern, pattern, pattern)
|
|
}
|
|
var response PageResponse
|
|
if err := query.Count(&response.Count).Error; err != nil {
|
|
return PageResponse{}, fmt.Errorf("count outbox messages: %w", err)
|
|
}
|
|
if err := query.Order("created_at DESC, id DESC").Limit(request.GetPageSize()).Offset((request.GetPageIndex() - 1) * request.GetPageSize()).Find(&response.List).Error; err != nil {
|
|
return PageResponse{}, fmt.Errorf("list outbox messages: %w", err)
|
|
}
|
|
for state, target := range map[string]*int64{StatePending: &response.Summary.Pending, StateRetry: &response.Summary.Retry, StateProcessing: &response.Summary.Processing, StateDead: &response.Summary.Dead} {
|
|
if err := s.Orm.Model(&Message{}).Where("state = ?", state).Count(target).Error; err != nil {
|
|
return PageResponse{}, fmt.Errorf("summarize outbox: %w", err)
|
|
}
|
|
}
|
|
return response, nil
|
|
}
|
|
|
|
func (s *Service) Get(id string) (DetailResponse, error) {
|
|
id = strings.TrimSpace(id)
|
|
if id == "" || len(id) > 64 {
|
|
return DetailResponse{}, ErrInvalidInput
|
|
}
|
|
var message Message
|
|
if err := s.Orm.First(&message, "id = ?", id).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return DetailResponse{}, ErrNotFound
|
|
}
|
|
return DetailResponse{}, fmt.Errorf("get outbox message: %w", err)
|
|
}
|
|
var attempts []Attempt
|
|
if err := s.Orm.Where("message_id = ?", id).Order("created_at DESC, id DESC").Limit(50).Find(&attempts).Error; err != nil {
|
|
return DetailResponse{}, fmt.Errorf("list outbox attempts: %w", err)
|
|
}
|
|
return DetailResponse{Message: message, Attempts: attempts}, nil
|
|
}
|
|
|
|
func (s *Service) Requeue(id string, request RequeueRequest, actorUserID int) (Message, error) {
|
|
reason := strings.TrimSpace(request.Reason)
|
|
if utf8.RuneCountInString(reason) < 6 || utf8.RuneCountInString(reason) > 256 || request.ExpectedVersion < 1 {
|
|
return Message{}, ErrInvalidInput
|
|
}
|
|
now := s.now()
|
|
var result Message
|
|
err := s.Orm.Transaction(func(tx *gorm.DB) error {
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&result, "id = ?", strings.TrimSpace(id)).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return ErrNotFound
|
|
}
|
|
return err
|
|
}
|
|
if result.Version != request.ExpectedVersion {
|
|
return ErrVersionConflict
|
|
}
|
|
if result.State != StateDead {
|
|
return ErrNotDead
|
|
}
|
|
result.State, result.AvailableAt, result.LastError = StatePending, now, ""
|
|
result.LeaseOwner, result.LeaseUntil, result.Version = "", nil, result.Version+1
|
|
if err := tx.Save(&result).Error; err != nil {
|
|
return err
|
|
}
|
|
return tx.Create(&Attempt{MessageID: result.ID, Number: result.AttemptCount, Outcome: "manual_requeue", Detail: reason, ActorUserID: actorUserID, CreatedAt: now}).Error
|
|
})
|
|
if err != nil {
|
|
return Message{}, err
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func (s *Service) now() time.Time {
|
|
if s.Now != nil {
|
|
return s.Now().UTC()
|
|
}
|
|
return time.Now().UTC()
|
|
}
|
|
func validState(value string) bool {
|
|
return value == StatePending || value == StateProcessing || value == StateRetry || value == StateDead || value == StateDelivered
|
|
}
|