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) }