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