service.xpcool.com/internal/service/job/job.go

233 lines
6.9 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// 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 implementation not registered")
}
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, "query auto jobs")
}
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, "update auto job")
}
// 动态重启调度:移除旧 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, "count job log")
}
var list []entity.AutoJobLog
if err = m.Page(f.Page, f.Size).OrderDesc("id").Scan(&list); err != nil {
return nil, 0, gerror.Wrap(err, "query job log")
}
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, "load auto jobs failed: %v", err)
return
}
for _, j := range jobs {
s.addCron(ctx, j)
}
g.Log().Infof(ctx, "auto job scheduler started: %d jobs", 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, "job code %s has no handler, skip cron", 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, "add cron %s(%s) failed: %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("no handler for job %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, "update job %s status failed: %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, "insert job log failed: %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)
}
}