75 lines
2.2 KiB
Go
75 lines
2.2 KiB
Go
package service
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"log"
|
|
|
|
"gorm.io/gorm"
|
|
|
|
"oci-portal/internal/model"
|
|
)
|
|
|
|
// AI 探测任务由系统按渠道数量自动管理:每天 00:00 探测一次全部号池渠道
|
|
// (逐渠道使用该渠道自己的租户凭据,与号池一一对应)。
|
|
const (
|
|
aiProbeTaskName = "AI渠道探测"
|
|
aiProbeCron = "0 0 * * *"
|
|
)
|
|
|
|
// SyncAiProbeTask 按渠道数量同步 AI 探测任务:
|
|
// 渠道数 > 0 时确保任务存在且激活,归零时连同日志自动删除;
|
|
// 由渠道增删钩子与启动时调用。cron 与代码基线不一致时一并对齐(升级自动迁移)。
|
|
func (s *TaskService) SyncAiProbeTask(ctx context.Context) {
|
|
var count int64
|
|
if err := s.db.WithContext(ctx).Model(&model.AiChannel{}).Count(&count).Error; err != nil {
|
|
return
|
|
}
|
|
var task model.Task
|
|
err := s.db.WithContext(ctx).Where("type = ?", model.TaskTypeAiProbe).First(&task).Error
|
|
notFound := errors.Is(err, gorm.ErrRecordNotFound)
|
|
if err != nil && !notFound {
|
|
return
|
|
}
|
|
if !notFound && task.CronExpr != aiProbeCron {
|
|
s.alignAiProbeCron(ctx, &task)
|
|
}
|
|
if count > 0 {
|
|
s.activateAiProbe(ctx, &task, notFound)
|
|
return
|
|
}
|
|
if !notFound {
|
|
if err := s.removeTask(ctx, task.ID); err != nil {
|
|
log.Printf("sync ai probe task: %v", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// alignAiProbeCron 把存量任务的执行计划对齐到当前基线。
|
|
func (s *TaskService) alignAiProbeCron(ctx context.Context, task *model.Task) {
|
|
cron := aiProbeCron
|
|
if _, err := s.UpdateTask(ctx, task.ID, UpdateTaskInput{CronExpr: &cron}); err != nil {
|
|
log.Printf("align ai probe cron: %v", err)
|
|
return
|
|
}
|
|
task.CronExpr = aiProbeCron
|
|
}
|
|
|
|
// activateAiProbe 创建或激活 AI 探测任务。
|
|
func (s *TaskService) activateAiProbe(ctx context.Context, task *model.Task, notFound bool) {
|
|
if notFound {
|
|
in := CreateTaskInput{Name: aiProbeTaskName, Type: model.TaskTypeAiProbe, CronExpr: aiProbeCron}
|
|
if _, err := s.CreateTask(ctx, in); err != nil {
|
|
log.Printf("sync ai probe task: %v", err)
|
|
}
|
|
return
|
|
}
|
|
if task.Status == model.TaskStatusActive {
|
|
return
|
|
}
|
|
active := model.TaskStatusActive
|
|
if _, err := s.UpdateTask(ctx, task.ID, UpdateTaskInput{Status: &active}); err != nil {
|
|
log.Printf("sync ai probe task: %v", err)
|
|
}
|
|
}
|