Files

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