Some checks failed
Build and Deploy (service.xpcool.com) / build-and-deploy (push) Failing after 31s
将 api 层校验规则提示语、service 层 gerror.Wrap 与 response.Error 错误信息、panic 未注册提示统一改为中文,并同步中文化 cmd 路由注释 与变更记录。仅涉及注释、文档与字符串改动,无业务逻辑变更。
233 lines
7.0 KiB
Go
233 lines
7.0 KiB
Go
// Package job 提供自动任务调度服务:DB 驱动的 gcron 管理 + 运行记录 + 完成/失败通知。
|
||
package job
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/gogf/gf/v2/errors/gerror"
|
||
"github.com/gogf/gf/v2/frame/g"
|
||
"github.com/gogf/gf/v2/os/gcron"
|
||
"github.com/gogf/gf/v2/os/gtime"
|
||
|
||
"service.xpcool.com/internal/consts"
|
||
"service.xpcool.com/internal/dao"
|
||
"service.xpcool.com/internal/library/response"
|
||
"service.xpcool.com/internal/model/dto"
|
||
"service.xpcool.com/internal/model/entity"
|
||
"service.xpcool.com/internal/service/notice"
|
||
)
|
||
|
||
// TaskFunc 任务执行函数:返回运行摘要与错误。
|
||
type TaskFunc func(context.Context) (string, error)
|
||
|
||
// IJob 自动任务服务接口。
|
||
type IJob interface {
|
||
List(context.Context) ([]dto.JobVO, error)
|
||
Save(context.Context, dto.JobInput) error
|
||
Trigger(context.Context, uint64) (string, bool, error)
|
||
LogList(context.Context, dto.JobLogFilter) ([]dto.JobLogVO, int, error)
|
||
Register(string, TaskFunc)
|
||
StartScheduler(context.Context)
|
||
}
|
||
|
||
type job struct {
|
||
mu sync.Mutex
|
||
registry map[string]TaskFunc // code -> 业务函数
|
||
}
|
||
|
||
var localJob IJob
|
||
|
||
// New 创建自动任务服务实现。
|
||
func New() IJob {
|
||
return &job{registry: make(map[string]TaskFunc)}
|
||
}
|
||
|
||
// Job 返回已注册的自动任务服务实现。
|
||
func Job() IJob {
|
||
if localJob == nil {
|
||
panic("Job 实现未注册")
|
||
}
|
||
return localJob
|
||
}
|
||
|
||
// RegisterJob 注册自动任务服务实现。
|
||
func RegisterJob(i IJob) { localJob = i }
|
||
|
||
// Register 注册任务执行函数(按 code)。
|
||
func (s *job) Register(code string, fn TaskFunc) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
s.registry[code] = fn
|
||
}
|
||
|
||
// ---------------- 列表 / 保存 / 触发 / 日志 ----------------
|
||
|
||
func (s *job) List(ctx context.Context) ([]dto.JobVO, error) {
|
||
var list []entity.AutoJob
|
||
if err := dao.AutoJob.Ctx(ctx).OrderAsc("id").Scan(&list); err != nil {
|
||
return nil, gerror.Wrap(err, "查询自动任务列表失败")
|
||
}
|
||
out := make([]dto.JobVO, 0, len(list))
|
||
for _, v := range list {
|
||
out = append(out, dto.JobVO{
|
||
Id: v.Id, Name: v.Name, Code: v.Code, JobType: v.JobType, CronExpr: v.CronExpr,
|
||
Enabled: v.Enabled, Remark: v.Remark, LastRunAt: v.LastRunAt,
|
||
LastResult: v.LastResult, LastError: v.LastError,
|
||
})
|
||
}
|
||
return out, nil
|
||
}
|
||
|
||
func (s *job) Save(ctx context.Context, in dto.JobInput) error {
|
||
var j entity.AutoJob
|
||
if err := dao.AutoJob.Ctx(ctx).Where("id", in.Id).Scan(&j); err != nil || j.Id == 0 {
|
||
return response.Error(consts.CodeInvalidParam, "任务不存在")
|
||
}
|
||
data := map[string]interface{}{"remark": in.Remark}
|
||
if in.CronExpr != "" {
|
||
data["cron_expr"] = in.CronExpr
|
||
}
|
||
if in.Enabled >= 0 {
|
||
data["enabled"] = in.Enabled
|
||
}
|
||
if _, err := dao.AutoJob.Ctx(ctx).Where("id", in.Id).Data(data).Update(); err != nil {
|
||
return gerror.Wrap(err, "更新自动任务失败")
|
||
}
|
||
// 动态重启调度:移除旧 cron,按新配置重新注册
|
||
s.reschedule(ctx, j.Code)
|
||
return nil
|
||
}
|
||
|
||
func (s *job) Trigger(ctx context.Context, id uint64) (string, bool, error) {
|
||
var j entity.AutoJob
|
||
if err := dao.AutoJob.Ctx(ctx).Where("id", id).Scan(&j); err != nil || j.Id == 0 {
|
||
return "", false, response.Error(consts.CodeInvalidParam, "任务不存在")
|
||
}
|
||
summary, err := s.run(ctx, j)
|
||
return summary, err == nil, err
|
||
}
|
||
|
||
func (s *job) LogList(ctx context.Context, f dto.JobLogFilter) ([]dto.JobLogVO, int, error) {
|
||
m := dao.AutoJobLog.Ctx(ctx)
|
||
if f.JobId > 0 {
|
||
m = m.Where("job_id", f.JobId)
|
||
}
|
||
total, err := m.Clone().Count()
|
||
if err != nil {
|
||
return nil, 0, gerror.Wrap(err, "统计任务日志总数失败")
|
||
}
|
||
var list []entity.AutoJobLog
|
||
if err = m.Page(f.Page, f.Size).OrderDesc("id").Scan(&list); err != nil {
|
||
return nil, 0, gerror.Wrap(err, "查询任务日志列表失败")
|
||
}
|
||
out := make([]dto.JobLogVO, 0, len(list))
|
||
for _, v := range list {
|
||
jobName := ""
|
||
if f.JobId == 0 {
|
||
var j entity.AutoJob
|
||
_ = dao.AutoJob.Ctx(ctx).Where("id", v.JobId).Scan(&j)
|
||
jobName = j.Name
|
||
}
|
||
out = append(out, dto.JobLogVO{
|
||
Id: v.Id, JobId: v.JobId, JobName: jobName, RunAt: v.RunAt,
|
||
Result: v.Result, Error: v.Error, Summary: v.Summary, DurationMs: v.DurationMs,
|
||
})
|
||
}
|
||
return out, total, nil
|
||
}
|
||
|
||
// ---------------- 调度 ----------------
|
||
|
||
// StartScheduler 启动调度:加载 DB 中启用的任务并注册 gcron。
|
||
func (s *job) StartScheduler(ctx context.Context) {
|
||
var jobs []entity.AutoJob
|
||
if err := dao.AutoJob.Ctx(ctx).Where("enabled", 1).Scan(&jobs); err != nil {
|
||
g.Log().Errorf(ctx, "加载自动任务失败: %v", err)
|
||
return
|
||
}
|
||
for _, j := range jobs {
|
||
s.addCron(ctx, j)
|
||
}
|
||
g.Log().Infof(ctx, "自动任务调度器已启动,共 %d 个任务", len(jobs))
|
||
}
|
||
|
||
// addCron 为任务注册 gcron(不存在的函数跳过)。
|
||
func (s *job) addCron(ctx context.Context, j entity.AutoJob) {
|
||
s.mu.Lock()
|
||
_, ok := s.registry[j.Code]
|
||
s.mu.Unlock()
|
||
if !ok {
|
||
g.Log().Warningf(ctx, "任务 %s 未注册执行函数,跳过定时注册", j.Code)
|
||
return
|
||
}
|
||
if _, err := gcron.Add(ctx, j.CronExpr, func(c context.Context) {
|
||
_, _ = s.run(c, j)
|
||
}, j.Code); err != nil {
|
||
g.Log().Errorf(ctx, "注册定时任务 %s(%s) 失败: %v", j.Code, j.CronExpr, err)
|
||
}
|
||
}
|
||
|
||
// reschedule 按 DB 最新状态重启任务的调度(停用则移除 cron)。
|
||
func (s *job) reschedule(ctx context.Context, code string) {
|
||
gcron.Remove(code)
|
||
var j entity.AutoJob
|
||
if err := dao.AutoJob.Ctx(ctx).Where("code", code).Scan(&j); err != nil || j.Id == 0 {
|
||
return
|
||
}
|
||
if j.Enabled == 1 {
|
||
s.addCron(ctx, j)
|
||
}
|
||
}
|
||
|
||
// run 执行任务:调用注册函数 → 更新任务状态 → 写运行日志 → 通知完成/失败。
|
||
func (s *job) run(ctx context.Context, j entity.AutoJob) (string, error) {
|
||
s.mu.Lock()
|
||
fn, ok := s.registry[j.Code]
|
||
s.mu.Unlock()
|
||
if !ok {
|
||
err := fmt.Errorf("任务 %s 未注册执行函数", j.Code)
|
||
s.finish(ctx, j, 2, err.Error(), "", 0)
|
||
return "", err
|
||
}
|
||
start := time.Now()
|
||
summary, err := fn(ctx)
|
||
duration := int(time.Since(start).Milliseconds())
|
||
if err != nil {
|
||
s.finish(ctx, j, 2, err.Error(), summary, duration)
|
||
return summary, err
|
||
}
|
||
s.finish(ctx, j, 1, "", summary, duration)
|
||
return summary, nil
|
||
}
|
||
|
||
// finish 更新任务状态 + 写运行日志 + 发通知。
|
||
func (s *job) finish(ctx context.Context, j entity.AutoJob, result int, errMsg, summary string, duration int) {
|
||
// 更新任务状态
|
||
if _, err := dao.AutoJob.Ctx(ctx).Where("id", j.Id).Data(map[string]interface{}{
|
||
"last_run_at": gtime.Now().Format("Y-m-d H:i:s"),
|
||
"last_result": result, "last_error": errMsg,
|
||
}).Update(); err != nil {
|
||
g.Log().Errorf(ctx, "更新任务 %s 状态失败: %v", j.Code, err)
|
||
}
|
||
// 写运行日志
|
||
if _, err := dao.AutoJobLog.Ctx(ctx).Data(map[string]interface{}{
|
||
"job_id": j.Id, "result": result, "error": errMsg,
|
||
"summary": summary, "duration_ms": duration,
|
||
}).Insert(); err != nil {
|
||
g.Log().Errorf(ctx, "写入任务运行日志失败: %v", err)
|
||
}
|
||
// 通知:完成/失败
|
||
vars := map[string]string{"jobName": j.Name, "summary": summary, "error": errMsg}
|
||
if result == 1 {
|
||
if summary == "" {
|
||
summary = "执行完成"
|
||
}
|
||
_ = notice.Notice().Send(ctx, "job_done", vars)
|
||
} else {
|
||
_ = notice.Notice().Send(ctx, "job_fail", vars)
|
||
}
|
||
}
|