package outbox import ( "context" stderrors "errors" "math" "sync" "time" "github.com/bytedance/sonic" "github.com/hibiken/asynq" "go.uber.org/zap" "gorm.io/gorm" "gorm.io/gorm/clause" "github.com/break/junhong_cmp_fiber/internal/model" "github.com/break/junhong_cmp_fiber/pkg/constants" "github.com/break/junhong_cmp_fiber/pkg/queue" ) // DeliveryEnvelope 是 Relay 原样传播到 Asynq 的公共结构化信封。 type DeliveryEnvelope struct { EventID string `json:"event_id"` EventType string `json:"event_type"` PayloadVersion int `json:"payload_version"` AggregateType string `json:"aggregate_type"` AggregateID string `json:"aggregate_id"` ResourceType string `json:"resource_type"` ResourceID string `json:"resource_id"` BusinessKey string `json:"business_key,omitempty"` RequestID string `json:"request_id,omitempty"` CorrelationID string `json:"correlation_id,omitempty"` Payload sonic.NoCopyRawMessage `json:"payload"` } // Publisher 是 Relay 的队列边界。 type Publisher interface { Publish(ctx context.Context, envelope DeliveryEnvelope) error } // PermanentError 表示重试无法修复的投递错误。 type PermanentError struct { err error } // Error 返回安全的错误文本,仅供内部日志和错误链判断使用。 func (e *PermanentError) Error() string { return e.err.Error() } // Unwrap 返回原始错误。 func (e *PermanentError) Unwrap() error { return e.err } // Permanent 将不可恢复错误标记为永久失败,Relay 会直接保留最终失败事实。 func Permanent(err error) error { if err == nil { return nil } return &PermanentError{err: err} } // QueuePublisher 使用项目统一队列客户端发布结构化信封。 type QueuePublisher struct { client *queue.Client } // NewQueuePublisher 创建项目统一队列客户端 Adapter。 func NewQueuePublisher(client *queue.Client) *QueuePublisher { return &QueuePublisher{client: client} } // Publish 将公共信封作为 struct 入队,禁止调用方预序列化。 func (p *QueuePublisher) Publish(ctx context.Context, envelope DeliveryEnvelope) error { return p.client.EnqueueTask(ctx, constants.TaskTypeOutboxDeliver, envelope) } // RelayOptions 控制单个 Relay 实例的领取与重试行为。 type RelayOptions struct { Owner string BatchSize int LeaseDuration time.Duration Now func() time.Time } // Relay 完成公共 Outbox 的租约领取和至少一次队列投递。 type Relay struct { db *gorm.DB publisher Publisher logger *zap.Logger options RelayOptions } // NewRelay 创建公共 Outbox Relay。 func NewRelay(db *gorm.DB, publisher Publisher, logger *zap.Logger, options RelayOptions) (*Relay, error) { if db == nil || publisher == nil || options.Owner == "" { return nil, stderrors.New("Outbox Relay 依赖和租约所有者不能为空") } if options.BatchSize <= 0 { options.BatchSize = constants.OutboxDefaultBatchSize } if options.LeaseDuration <= 0 { options.LeaseDuration = constants.OutboxDefaultLeaseDuration } if options.Now == nil { options.Now = time.Now } if logger == nil { logger = zap.NewNop() } return &Relay{db: db, publisher: publisher, logger: logger, options: options}, nil } // ProcessBatch 领取并投递一批到期事件。 func (r *Relay) ProcessBatch(ctx context.Context) (int, error) { events, err := r.ClaimBatch(ctx) if err != nil { return 0, err } processed := 0 for _, event := range events { if err := r.publisher.Publish(ctx, deliveryEnvelope(event)); err != nil { var permanentError *PermanentError permanent := stderrors.As(err, &permanentError) code := "OUTBOX_ENQUEUE_FAILED" summary := "队列暂时不可用" if permanent { code = "OUTBOX_PERMANENT_FAILURE" summary = "事件无法投递,已停止自动重试" } if failErr := r.markFailed(ctx, event, code, summary, permanent); failErr != nil { return processed, failErr } r.logger.Warn("Outbox 事件投递失败", zap.String("event_id", event.EventID), zap.String("correlation_id", event.CorrelationID), zap.String("error_code", code), zap.Bool("permanent", permanent)) continue } if err := r.MarkDelivered(ctx, event.ID); err != nil { // 入队成功但标记失败时保留租约,过期后会使用同一 event_id 再次投递。 return processed, err } processed++ } return processed, nil } // ClaimBatch 通过行锁跳过竞争行,并领取待投递或租约过期事件。 func (r *Relay) ClaimBatch(ctx context.Context) ([]model.OutboxEvent, error) { now := r.options.Now().UTC() expiresAt := now.Add(r.options.LeaseDuration) claimed := make([]model.OutboxEvent, 0, r.options.BatchSize) err := r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { var candidates []model.OutboxEvent if err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}). Where("(status = ? AND next_attempt_at <= ?) OR (status = ? AND lease_expires_at <= ?)", constants.OutboxStatusPending, now, constants.OutboxStatusDelivering, now). Order("created_at ASC, id ASC").Limit(r.options.BatchSize).Find(&candidates).Error; err != nil { return err } for index := range candidates { result := tx.Model(&model.OutboxEvent{}). Where("id = ? AND ((status = ? AND next_attempt_at <= ?) OR (status = ? AND lease_expires_at <= ?))", candidates[index].ID, constants.OutboxStatusPending, now, constants.OutboxStatusDelivering, now). Updates(map[string]any{ "status": constants.OutboxStatusDelivering, "lease_owner": r.options.Owner, "lease_expires_at": expiresAt, "updated_at": now, }) if result.Error != nil { return result.Error } if result.RowsAffected == 1 { candidates[index].Status = constants.OutboxStatusDelivering candidates[index].LeaseOwner = &r.options.Owner candidates[index].LeaseExpiresAt = &expiresAt claimed = append(claimed, candidates[index]) } } return nil }) return claimed, err } // RenewLease 仅允许当前租约所有者续租投递中的事件。 func (r *Relay) RenewLease(ctx context.Context, eventID uint) (bool, error) { now := r.options.Now().UTC() result := r.db.WithContext(ctx).Model(&model.OutboxEvent{}). Where("id = ? AND status = ? AND lease_owner = ? AND lease_expires_at > ?", eventID, constants.OutboxStatusDelivering, r.options.Owner, now). Updates(map[string]any{"lease_expires_at": now.Add(r.options.LeaseDuration), "updated_at": now}) return result.RowsAffected == 1, result.Error } // MarkDelivered 仅允许当前租约所有者把事件标记为已投递。 func (r *Relay) MarkDelivered(ctx context.Context, eventID uint) error { now := r.options.Now().UTC() result := r.db.WithContext(ctx).Model(&model.OutboxEvent{}). Where("id = ? AND status = ? AND lease_owner = ?", eventID, constants.OutboxStatusDelivering, r.options.Owner). Updates(map[string]any{ "status": constants.OutboxStatusDelivered, "delivered_at": now, "lease_owner": nil, "lease_expires_at": nil, "last_error_code": "", "last_error_summary": "", "updated_at": now, }) if result.Error != nil { return result.Error } if result.RowsAffected != 1 { return stderrors.New("Outbox 投递完成条件不满足") } return nil } // MarkFailed 记录安全错误并退避;达到上限后保留最终失败事实。 func (r *Relay) MarkFailed(ctx context.Context, event model.OutboxEvent, code, summary string) error { return r.markFailed(ctx, event, code, summary, false) } func (r *Relay) markFailed(ctx context.Context, event model.OutboxEvent, code, summary string, permanent bool) error { now := r.options.Now().UTC() retryCount := event.RetryCount + 1 status := constants.OutboxStatusPending nextAttemptAt := now.Add(backoff(retryCount)) if permanent || retryCount >= event.MaxRetries { status = constants.OutboxStatusFailed nextAttemptAt = now r.logger.Error("Outbox 事件停止自动重试", zap.String("event_id", event.EventID), zap.String("correlation_id", event.CorrelationID), zap.String("error_code", code), zap.Bool("permanent", permanent)) } result := r.db.WithContext(ctx).Model(&model.OutboxEvent{}). Where("id = ? AND status = ? AND lease_owner = ?", event.ID, constants.OutboxStatusDelivering, r.options.Owner). Updates(map[string]any{ "status": status, "retry_count": retryCount, "next_attempt_at": nextAttemptAt, "lease_owner": nil, "lease_expires_at": nil, "last_error_code": code, "last_error_summary": summary, "updated_at": now, }) if result.Error != nil { return result.Error } if result.RowsAffected != 1 { return stderrors.New("Outbox 失败更新条件不满足") } return nil } func deliveryEnvelope(event model.OutboxEvent) DeliveryEnvelope { return DeliveryEnvelope{ EventID: event.EventID, EventType: event.EventType, PayloadVersion: event.PayloadVersion, AggregateType: event.AggregateType, AggregateID: event.AggregateID, ResourceType: event.ResourceType, ResourceID: event.ResourceID, BusinessKey: event.BusinessKey, RequestID: event.RequestID, CorrelationID: event.CorrelationID, Payload: sonic.NoCopyRawMessage(event.Payload), } } func backoff(retryCount int) time.Duration { delay := float64(constants.OutboxBaseRetryDelay) * math.Pow(2, float64(retryCount-1)) if delay > float64(constants.OutboxMaxRetryDelay) { return constants.OutboxMaxRetryDelay } return time.Duration(delay) } // EventConsumer 是业务消费者公开实现的事件处理边界。 type EventConsumer interface { Consume(ctx context.Context, envelope DeliveryEnvelope) error } // ConsumerRegistry 按稳定事件类型分发到业务消费者。 type ConsumerRegistry struct { mu sync.RWMutex consumers map[string]EventConsumer } // NewConsumerRegistry 创建空消费者注册表。 func NewConsumerRegistry() *ConsumerRegistry { return &ConsumerRegistry{consumers: map[string]EventConsumer{}} } // Register 注册一个由业务 PRD 拥有的事件消费者。 func (r *ConsumerRegistry) Register(eventType string, consumer EventConsumer) error { if eventType == "" || consumer == nil { return stderrors.New("Outbox 消费者注册信息不完整") } r.mu.Lock() defer r.mu.Unlock() if _, exists := r.consumers[eventType]; exists { return stderrors.New("Outbox 事件类型重复注册") } r.consumers[eventType] = consumer return nil } // Consume 将公共信封交给对应业务消费者;公共层不实现业务副作用。 func (r *ConsumerRegistry) Consume(ctx context.Context, envelope DeliveryEnvelope) error { r.mu.RLock() consumer := r.consumers[envelope.EventType] r.mu.RUnlock() if consumer == nil { return stderrors.New("Outbox 事件消费者尚未注册") } return consumer.Consume(ctx, envelope) } // Handler 是公共 Outbox Asynq 任务 Handler。 type Handler struct { consumer EventConsumer } // NewHandler 创建公共 Outbox Asynq Handler。 func NewHandler(consumer EventConsumer) *Handler { return &Handler{consumer: consumer} } // Handle 解析结构化信封并调用公开消费者边界。 func (h *Handler) Handle(ctx context.Context, task *asynq.Task) error { var envelope DeliveryEnvelope if err := sonic.Unmarshal(task.Payload(), &envelope); err != nil { return err } if envelope.EventID == "" || envelope.EventType == "" || envelope.PayloadVersion <= 0 { return stderrors.New("Outbox 事件信封不完整") } return h.consumer.Consume(ctx, envelope) }