service.xpcool.com/internal/service/job/job.go
夏犀麟 4b3c00270f
Some checks failed
Build and Deploy (service.xpcool.com) / build-and-deploy (push) Failing after 31s
feat(i18n): 中文化校验提示与服务层错误信息
将 api 层校验规则提示语、service 层 gerror.Wrap 与 response.Error
错误信息、panic 未注册提示统一改为中文,并同步中文化 cmd 路由注释
与变更记录。仅涉及注释、文档与字符串改动,无业务逻辑变更。
2026-09-13 23:22:46 +08:00

233 lines
7.0 KiB
Go
Raw Permalink 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 实现未注册")
}
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)
}
}