Files

218 lines
5.6 KiB
Go

package service
import (
"context"
"crypto/rand"
"encoding/hex"
"encoding/json"
"errors"
"log/slog"
"time"
"silk-server-go/internal/model"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
const (
OutboxStatusPending = "pending"
OutboxStatusSending = "sending"
OutboxStatusSent = "sent"
OutboxStatusRetry = "retry"
OutboxStatusFailed = "failed"
OutboxStatusCancelled = "cancelled"
OutboxEventWechatSubscribe = "wechat.subscribe"
OutboxEventDetectionTaskCreate = "detection.task.create"
)
// Event 写入 outbox 的领域事件。
type Event struct {
ID string
Type string
AggregateType string
AggregateID string
Payload json.RawMessage
}
// EventHandler 处理一条已领取的 outbox 事件。
type EventHandler func(ctx context.Context, event model.OutboxEvent) error
// Outbox 基于 PostgreSQL 的可靠事件箱。
type Outbox struct {
db *gorm.DB
handler EventHandler
interval time.Duration
batchSize int
maxAttempts int
backoff time.Duration
}
// NewOutbox 创建默认配置的 Outbox。
func NewOutbox(db *gorm.DB) *Outbox {
return &Outbox{
db: db,
interval: 5 * time.Second,
batchSize: 20,
maxAttempts: 5,
backoff: 15 * time.Second,
}
}
// SetHandler 设置事件处理器。
func (o *Outbox) SetHandler(handler EventHandler) {
o.handler = handler
}
// PublishTx 在业务事务内写入事件;重复 eventId 通过唯一索引幂等跳过。
func (o *Outbox) PublishTx(tx *gorm.DB, event Event) error {
if event.Type == "" {
return errors.New("outbox event type is required")
}
if event.ID == "" {
event.ID = randomEventID()
}
payload := event.Payload
if len(payload) == 0 {
payload = json.RawMessage(`{}`)
}
record := model.OutboxEvent{
EventID: event.ID,
EventType: event.Type,
AggregateType: strPtrOrNil(event.AggregateType),
AggregateID: strPtrOrNil(event.AggregateID),
Payload: payload,
Status: OutboxStatusPending,
MaxAttempts: o.maxAttempts,
NextAttemptAt: time.Now(),
}
// PublishTx 必须复用调用方事务,避免 GORM 对单条 Create 再开嵌套事务。
return tx.Session(&gorm.Session{SkipDefaultTransaction: true}).Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "event_id"}},
DoNothing: true,
}).Create(&record).Error
}
// ClaimNext 领取一批到期待处理事件。
func (o *Outbox) ClaimNext(ctx context.Context, limit int) ([]model.OutboxEvent, error) {
if limit <= 0 {
limit = o.batchSize
}
var events []model.OutboxEvent
err := o.db.WithContext(ctx).
Where("status IN ? AND next_attempt_at <= ?", []string{OutboxStatusPending, OutboxStatusRetry}, time.Now()).
Order("created_at ASC").
Limit(limit).
Find(&events).Error
return events, err
}
// MarkSending 原子标记领取状态,避免多 worker 重复处理。
func (o *Outbox) MarkSending(ctx context.Context, id string) (bool, error) {
result := o.db.WithContext(ctx).
Model(&model.OutboxEvent{}).
Where("id = ? AND status IN ?", id, []string{OutboxStatusPending, OutboxStatusRetry}).
Update("status", OutboxStatusSending)
return result.RowsAffected > 0, result.Error
}
// MarkResult 写入成功/失败并计算下一次重试时间。
func (o *Outbox) MarkResult(ctx context.Context, id string, processErr error) error {
var event model.OutboxEvent
if err := o.db.WithContext(ctx).First(&event, "id = ?", id).Error; err != nil {
return err
}
if event.MaxAttempts <= 0 {
event.MaxAttempts = o.maxAttempts
}
event.Attempts++
now := time.Now()
updates := map[string]interface{}{
"attempts": event.Attempts,
"updated_at": now,
"next_attempt_at": now,
}
if processErr == nil {
updates["status"] = OutboxStatusSent
updates["last_error"] = nil
} else {
message := processErr.Error()
updates["last_error"] = message
if event.Attempts >= event.MaxAttempts {
updates["status"] = OutboxStatusFailed
} else {
updates["status"] = OutboxStatusRetry
updates["next_attempt_at"] = now.Add(o.backoff * time.Duration(event.Attempts))
}
}
return o.db.WithContext(ctx).
Model(&model.OutboxEvent{}).
Where("id = ?", id).
Updates(updates).Error
}
// ProcessPending 领取并处理到期事件,返回本轮处理数量。
func (o *Outbox) ProcessPending(ctx context.Context) (int, error) {
if o.handler == nil {
return 0, errors.New("outbox handler is not set")
}
processed := 0
for {
events, err := o.ClaimNext(ctx, o.batchSize)
if err != nil {
return processed, err
}
if len(events) == 0 {
return processed, nil
}
for _, event := range events {
claimed, err := o.MarkSending(ctx, event.ID)
if err != nil {
slog.Warn("outbox mark sending failed", "eventId", event.EventID, "err", err)
continue
}
if !claimed {
continue
}
processErr := o.handler(ctx, event)
if err := o.MarkResult(ctx, event.ID, processErr); err != nil {
slog.Warn("outbox mark result failed", "eventId", event.EventID, "err", err)
continue
}
processed++
}
}
}
// Start 启动后台 worker。
func (o *Outbox) Start(ctx context.Context) {
ticker := time.NewTicker(o.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if _, err := o.ProcessPending(ctx); err != nil {
slog.Warn("outbox worker failed", "err", err)
}
}
}
}
func strPtrOrNil(value string) *string {
if value == "" {
return nil
}
return &value
}
func randomEventID() string {
b := make([]byte, 16)
if _, err := rand.Read(b); err != nil {
return hex.EncodeToString([]byte(time.Now().Format(time.RFC3339Nano)))
}
return hex.EncodeToString(b)
}