package task import ( "context" "fmt" "strings" "time" "github.com/bytedance/sonic" "github.com/hibiken/asynq" "github.com/redis/go-redis/v9" "go.uber.org/zap" "github.com/break/junhong_cmp_fiber/pkg/constants" ) // EmailPayload 邮件任务载荷 type EmailPayload struct { RequestID string `json:"request_id"` To string `json:"to"` Subject string `json:"subject"` Body string `json:"body"` CC []string `json:"cc,omitempty"` Attachments []string `json:"attachments,omitempty"` } // EmailHandler 邮件任务处理器 type EmailHandler struct { redis *redis.Client logger *zap.Logger } // NewEmailHandler 创建邮件任务处理器 func NewEmailHandler(redis *redis.Client, logger *zap.Logger) *EmailHandler { return &EmailHandler{ redis: redis, logger: logger, } } // HandleEmailSend 处理邮件发送任务 func (h *EmailHandler) HandleEmailSend(ctx context.Context, task *asynq.Task) error { // 解析任务载荷 var payload EmailPayload if err := sonic.Unmarshal(task.Payload(), &payload); err != nil { h.logger.Error("解析邮件任务载荷失败", zap.Error(err), zap.String("task_id", task.ResultWriter().TaskID()), ) return asynq.SkipRetry // JSON 解析失败不重试 } // 验证载荷 if err := h.validatePayload(&payload); err != nil { h.logger.Error("邮件任务载荷验证失败", zap.Error(err), zap.String("request_id", payload.RequestID), ) return asynq.SkipRetry // 参数错误不重试 } // 幂等性检查:使用 Redis 锁 lockKey := constants.RedisTaskLockKey(payload.RequestID) locked, err := h.acquireLock(ctx, lockKey) if err != nil { h.logger.Error("获取任务锁失败", zap.Error(err), zap.String("request_id", payload.RequestID), ) return err // 锁获取失败,可以重试 } if !locked { h.logger.Info("任务已执行,跳过(幂等性)", zap.String("request_id", payload.RequestID), zap.String("to", payload.To), ) return nil // 已执行,跳过 } // 记录任务开始执行 h.logger.Info("开始处理邮件发送任务", zap.String("request_id", payload.RequestID), zap.String("to", payload.To), zap.String("subject", payload.Subject), zap.Int("cc_count", len(payload.CC)), zap.Int("attachments_count", len(payload.Attachments)), ) // 执行邮件发送(模拟) if err := h.sendEmail(ctx, &payload); err != nil { h.logger.Error("邮件发送失败", zap.Error(err), zap.String("request_id", payload.RequestID), zap.String("to", payload.To), ) return err // 发送失败,可以重试 } // 记录任务完成 h.logger.Info("邮件发送成功", zap.String("request_id", payload.RequestID), zap.String("to", payload.To), ) return nil } // validatePayload 验证邮件载荷 func (h *EmailHandler) validatePayload(payload *EmailPayload) error { if payload.RequestID == "" { return fmt.Errorf("request_id 不能为空") } if payload.To == "" { return fmt.Errorf("收件人不能为空") } if !strings.Contains(payload.To, "@") { return fmt.Errorf("邮箱格式无效") } if payload.Subject == "" { return fmt.Errorf("邮件主题不能为空") } if payload.Body == "" { return fmt.Errorf("邮件正文不能为空") } return nil } // acquireLock 获取 Redis 锁(幂等性) func (h *EmailHandler) acquireLock(ctx context.Context, key string) (bool, error) { // 使用 SetNX 实现分布式锁 // 过期时间 24 小时,防止锁永久存在 result, err := h.redis.SetNX(ctx, key, "1", 24*time.Hour).Result() if err != nil { return false, fmt.Errorf("设置 Redis 锁失败: %w", err) } return result, nil } // sendEmail 发送邮件(实际实现需要集成 SMTP 或邮件服务) func (h *EmailHandler) sendEmail(ctx context.Context, payload *EmailPayload) error { // TODO: 实际实现中需要集成邮件发送服务 // 例如:使用 SMTP、SendGrid、AWS SES 等 // 模拟发送延迟 time.Sleep(100 * time.Millisecond) // 这里仅作演示,实际应用中需要调用真实的邮件发送 API h.logger.Debug("模拟邮件发送", zap.String("to", payload.To), zap.String("subject", payload.Subject), zap.Int("body_length", len(payload.Body)), ) return nil }