* fix(database): correct name entry for SingBox in initialization script * fix(purchase): update gift amount deduction logic and handle zero-amount order status * feat: add type and default fields to rule group requests and update related logic * feat(rule): implement logic to set a default rule group during creation and update * fix(rule): add type and default fields to rule group model and update related logic * feat(proxy): enhance proxy group handling and sorting logic * refactor(proxy): replace hardcoded group names with constants for better maintainability * fix(proxy): update group selection logic to skip empty and default names * feat(proxy): enhance proxy and group handling with new configuration options * feat(surge): add Surge adapter support and enhance subscription URL handling * feat(traffic): implement traffic reset logic for subscription cycles * feat(auth): improve email and mobile config unmarshalling with default values * fix(auth) upbind email not update * fix(order) discount set default 1 * fix(order) discount set default 1 * fix: refactor surfboard proxy handling and enhance configuration template * fix(renewal) discount set default 1 * feat(loon): add Loon configuration template and enhance proxy handling * feat(subscription): update user subscription status based on expiration time * fix(renewal): update subscription retrieval method to use token instead of order ID * feat(order): enhance order processing logic with improved error handling and user subscription management * fix(order): improve code quality and fix critical bugs in order processing logic - Fix inconsistent logging calls across all order logic files - Fix critical gift amount deduction logic bug in renewal process - Fix variable shadowing errors in database transactions - Add comprehensive Go-standard documentation comments - Improve log prefix consistency for better debugging - Remove redundant discount validation code * fix(docker): add build argument for version in Docker image build process * feat(version): add endpoint to retrieve application version information * fix(auth): improve user authentication method logic and update user cache * feat(user): add ordering functionality to user list retrieval * fix(RevenueStatistics) fill list * fix(UserStatistics) fill list * fix(user): implement user cache clearing after auth method operations * fix(auth): enhance OAuth login logic with improved request handling and user registration flow * fix(user): implement sorting for authentication methods based on priority * fix(user): correct ordering clause for user retrieval based on filter * refactor(user): streamline cache management and enhance cache clearing logic * feat(logs) set logs volume in develop * fix(handler): implement browser interception to deny access for specific user agents * fix(resetTraffic) reset daily server * refactor(trojan): remove unused parameter and clean up logging in slice * fix(middleware): add domain length check and improve user-agent handling * fix(middleware): reorder domain processing and enhance user-agent handling * fix(resetTraffic): update subscription reset logic to use expire_time for monthly and yearly checks * fix(scheduler): update reset traffic task schedule to run daily at 00:30 * fix(traffic): enhance traffic reset logic for subscriptions and adjust status checks * fix(activateOrder): update traffic reset logic to include reset day check * feat(marketing): add batch email task management API and logic * feat(application): implement CRUD operations for subscribe applications * feat(types): add user agent limit and list to subscription configuration * feat(application): update subscription application requests to include structured download links * feat(application): add scheme field and download link handling to subscribe application * feat(application): add endpoint to retrieve client information * feat(application): move DownloadLink and SubscribeApplication types to types.api * feat(application): add DownloadLink and SubscribeClient types, update client response structure * feat(application): remove ProxyTemplate field from application API * feat(application): implement adapter for client configuration and add preview template functionality * feat(application): move DownloadLink type to types.api and remove from common.api * feat(application): update PreviewSubscribeTemplate to return structured response * feat(application): remove ProxyTemplate field from application API * feat(application): enhance cache key generation for user list and server data * feat(subscribe): add ClearCache method to manage subscription cache invalidation * feat(payment): add Description field to PaymentMethodDetail response * feat(subscribe): update next reset time calculation to use ExpireTime * feat(purchase): include handling fee in total amount calculation * feat(subscribe): add V2SubscribeHandler and logic for enhanced subscription management * feat(subscribe): add output format configuration to subscription adapter * feat(application): default data --------- Co-authored-by: Chang lue Tsen <tension@ppanel.dev> Co-authored-by: NoWay <Bob455668@hotmail.com>
165 lines
4.3 KiB
Go
165 lines
4.3 KiB
Go
package email
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/perfect-panel/server/internal/model/task"
|
||
"github.com/perfect-panel/server/pkg/logger"
|
||
"gorm.io/gorm"
|
||
)
|
||
|
||
type ErrorInfo struct {
|
||
Error string `json:"error"`
|
||
Email string `json:"email"`
|
||
Time int64 `json:"time"`
|
||
}
|
||
|
||
type Worker struct {
|
||
id int64 // 任务ID
|
||
db *gorm.DB // 数据库连接
|
||
ctx context.Context // 上下文
|
||
sender Sender // 邮件发送器接口
|
||
status uint8 // 任务状态,0 表示未运行,1 表示运行中 2 表示已完成
|
||
}
|
||
|
||
func NewWorker(ctx context.Context, id int64, db *gorm.DB, sender Sender) *Worker {
|
||
return &Worker{
|
||
id: id,
|
||
db: db,
|
||
ctx: ctx,
|
||
sender: sender,
|
||
}
|
||
}
|
||
|
||
// GetID 获取Worker的任务ID
|
||
func (w *Worker) GetID() int64 {
|
||
return w.id
|
||
}
|
||
|
||
// IsRunning 检查Worker是否正在运行
|
||
func (w *Worker) IsRunning() uint8 {
|
||
return w.status
|
||
}
|
||
|
||
// Start 启动Worker,开始处理任务
|
||
func (w *Worker) Start() {
|
||
// 检查并发限制
|
||
limit.Lock()
|
||
defer limit.Unlock()
|
||
tx := w.db.WithContext(w.ctx)
|
||
var taskInfo task.EmailTask
|
||
if err := tx.Model(&task.EmailTask{}).Where("id = ?", w.id).First(&taskInfo).Error; err != nil {
|
||
logger.Error("Batch Send Email",
|
||
logger.Field("message", "Failed to find task"),
|
||
logger.Field("error", err.Error()),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
w.status = 2 // 设置状态为已完成
|
||
return
|
||
}
|
||
if taskInfo.Status != 0 {
|
||
logger.Error("Batch Send Email",
|
||
logger.Field("message", "Task already completed or in progress"),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
w.status = 2 // 设置状态为已完成
|
||
return
|
||
}
|
||
if taskInfo.Recipients == "" && taskInfo.Additional == "" {
|
||
logger.Error("Batch Send Email",
|
||
logger.Field("message", "No recipients or additional emails provided"),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
w.status = 2 // 设置状态为已完成
|
||
return
|
||
}
|
||
w.status = 1 // 设置状态为运行中
|
||
var recipients []string
|
||
// 解析收件人
|
||
if taskInfo.Recipients != "" {
|
||
recipients = append(recipients, strings.Split(taskInfo.Recipients, "\n")...)
|
||
}
|
||
// 解析附加收件人
|
||
if taskInfo.Additional != "" {
|
||
recipients = append(recipients, strings.Split(taskInfo.Additional, "\n")...)
|
||
}
|
||
if len(recipients) == 0 {
|
||
logger.Error("Batch Send Email",
|
||
logger.Field("message", "No valid recipients found"),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
w.status = 2 // 设置状态为已完成
|
||
return
|
||
}
|
||
|
||
// 设置发送间隔时间
|
||
var intervalTime time.Duration
|
||
if taskInfo.Interval == 0 {
|
||
intervalTime = 1 * time.Second
|
||
} else {
|
||
intervalTime = time.Duration(taskInfo.Interval) * time.Second
|
||
}
|
||
|
||
var errors []ErrorInfo
|
||
var count uint64
|
||
for _, recipient := range recipients {
|
||
select {
|
||
case <-w.ctx.Done():
|
||
logger.Info("Batch Send Email",
|
||
logger.Field("message", "Worker stopped by context cancellation"),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
return
|
||
default:
|
||
}
|
||
|
||
if taskInfo.Status == 0 {
|
||
taskInfo.Status = 1 // 1 表示任务进行中
|
||
}
|
||
|
||
if err := w.sender.Send([]string{recipient}, taskInfo.Subject, taskInfo.Content); err != nil {
|
||
logger.Error("Batch Send Email",
|
||
logger.Field("message", "Failed to send email"),
|
||
logger.Field("error", err.Error()),
|
||
logger.Field("recipient", recipient),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
errors = append(errors, ErrorInfo{
|
||
Error: err.Error(),
|
||
Email: recipient,
|
||
Time: time.Now().Unix(),
|
||
})
|
||
text, _ := json.Marshal(errors)
|
||
taskInfo.Errors = string(text)
|
||
}
|
||
count++
|
||
if err := tx.Model(&task.EmailTask{}).Save(&taskInfo).Error; err != nil {
|
||
logger.Error("Batch Send Email",
|
||
logger.Field("message", "Failed to update task progress"),
|
||
logger.Field("error", err.Error()),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
w.status = 2 // 设置状态为已完成
|
||
return
|
||
}
|
||
time.Sleep(intervalTime)
|
||
}
|
||
taskInfo.Status = 2 // 设置状态为已完成
|
||
if err := tx.Model(&task.EmailTask{}).Save(&taskInfo).Error; err != nil {
|
||
logger.Error("Batch Send Email",
|
||
logger.Field("message", "Failed to finalize task"),
|
||
logger.Field("error", err.Error()),
|
||
logger.Field("task_id", w.id),
|
||
)
|
||
} else {
|
||
logger.Info("Batch Send Email",
|
||
logger.Field("message", "Task completed successfully"),
|
||
logger.Field("task_id", w.id),
|
||
logger.Field("total_sent", count),
|
||
)
|
||
}
|
||
}
|