提交 61cd9aa8 authored 作者: mooncake9527's avatar mooncake9527

update

上级 f591f912
...@@ -4,7 +4,7 @@ import ( ...@@ -4,7 +4,7 @@ import (
"context" "context"
"fmt" "fmt"
"os" "os"
"runtime" "runtime/debug"
"strconv" "strconv"
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
...@@ -63,8 +63,8 @@ func (x *Application) Run() { ...@@ -63,8 +63,8 @@ func (x *Application) Run() {
func Recover() { func Recover() {
if r := recover(); r != nil { if r := recover(); r != nil {
pc, file, line, _ := runtime.Caller(3) logger.Fatalf("Recovered panic: %v\nStack trace:\n%s",
logger.Fatalf("recovered panic at %s[%s:%d]: %v\n", runtime.FuncForPC(pc).Name(), file, line, r) r, string(debug.Stack()))
} }
} }
......
...@@ -31,17 +31,21 @@ func ParseConf() (err error) { ...@@ -31,17 +31,21 @@ func ParseConf() (err error) {
if confType == "local" { if confType == "local" {
localConf() localConf()
logger.Infof("[conf]conf:%s", Cfg.ConfCenter.Local.Conf) logger.Infof("[conf]conf:%s", Cfg.ConfCenter.Local.Conf)
return parseLocal(Cfg.ConfCenter.Local.Conf, Cfg, func() { if err := parseLocal(Cfg.ConfCenter.Local.Conf, Cfg, func() {
ParseExtend() ParseExtend()
}) }); err != nil {
return err
}
} }
if confType == "consul" { if confType == "consul" {
consulConf() consulConf()
confKey := Cfg.ConfCenter.Consul.Conf confKey := Cfg.ConfCenter.Consul.Conf
logger.Infof("[conf]conf:%s", confKey) logger.Infof("[conf]conf:%s", confKey)
return NewFetcher(confKey, Cfg, func(conf []byte) { if err := NewFetcher(confKey, Cfg, func(conf []byte) {
ParseExtend() ParseExtend()
}).Start() }).Start(); err != nil {
return err
}
} }
if confType == "nacos" { if confType == "nacos" {
nacosConf() nacosConf()
......
package controller
import (
"context"
"errors"
"fmt"
"time"
"github.com/gin-gonic/gin"
"github.com/go-redsync/redsync/v4"
"gitlab.wanzhuangkj.com/tush/xpkg/gin/response"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/gocron"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/logger"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/ctxUtils"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/redSyncUtils"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/xtime"
"gitlab.wanzhuangkj.com/tush/xpkg/xcommon"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/biz"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/dao"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/enums"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/models"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/service"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/types"
)
const (
PWD = "1937396767217688576"
)
type cronJobController struct {
*xcommon.Controller
}
var CronJobController cronJobController
func (x *cronJobController) StartByCode(c *gin.Context) {
req := &types.CronJobCodeReq{}
if err := c.ShouldBind(req); err != nil {
response.Error(c, err)
return
}
if req.Pwd != PWD {
response.Error(c, errors.New("pwd err"))
return
}
ctx := ctxUtils.WrapCtx(c)
job, err := service.CronJobService.GetByJobCode(ctx, req.JobCode, &biz.CronJobOpts{})
if err != nil {
response.Error(c, err)
return
}
if err := service.CronJobService.StartByID(ctx, job.ID, nil); err != nil {
response.Error(c, err)
return
}
gocron.Add(&gocron.Task{
TimeSpec: job.TimeSpec,
Name: job.JobCode,
Fn: func() {
ExecJob(context.TODO(), job)
},
})
response.Success(c)
}
func (x *cronJobController) StopByCode(c *gin.Context) {
req := &types.CronJobCodeReq{}
if err := c.ShouldBind(req); err != nil {
response.Error(c, err)
return
}
if req.Pwd != PWD {
response.Error(c, errors.New("pwd err"))
return
}
ctx := ctxUtils.WrapCtx(c)
job, err := service.CronJobService.GetByJobCode(ctx, req.JobCode, &biz.CronJobOpts{})
if err != nil {
response.Error(c, err)
return
}
if err := service.CronJobService.StopByID(ctx, job.ID, nil); err != nil {
response.Error(c, err)
return
}
gocron.DeleteTask(job.JobCode)
response.Success(c)
}
func (x *cronJobController) CallOnceByCode(c *gin.Context) {
req := &types.CronJobCodeReq{}
if err := c.ShouldBind(req); err != nil {
response.Error(c, err)
return
}
if req.Pwd != PWD {
response.Error(c, errors.New("pwd err"))
return
}
job, err := service.CronJobService.GetByJobCode(ctxUtils.WrapCtx(c), req.JobCode, &biz.CronJobOpts{})
if err != nil {
response.Error(c, err)
return
}
err = ExecJob(ctxUtils.WrapCtx(c), job)
if err != nil {
response.Error(c, err)
return
}
response.Success(c)
}
func (x *cronJobController) Create(c *gin.Context) {
in := &types.CronJobCreateReq{}
if err := x.Bind(c, in); err != nil {
response.Error(c, err)
return
}
id, err := service.CronJobService.Create(ctxUtils.WrapCtx(c), in, nil)
if err != nil {
response.Error(c, err)
return
}
response.Success(c, id)
}
func (x *cronJobController) DeleteByID(c *gin.Context) {
in := &types.BaseIDReq{}
if err := x.Bind(c, in); err != nil {
response.Error(c, err)
return
}
if err := service.CronJobService.DeleteByID(ctxUtils.WrapCtx(c), in.ID, nil); err != nil {
response.Error(c, err)
return
}
response.Success(c)
}
func (x *cronJobController) UpdateByID(c *gin.Context) {
in := &types.CronJobUpdateByIDReq{}
if err := x.Bind(c, in); err != nil {
response.Error(c, err)
return
}
if err := service.CronJobService.UpdateByID(ctxUtils.WrapCtx(c), in, nil); err != nil {
response.Error(c, err)
return
}
response.Success(c)
}
func (x *cronJobController) GetByID(c *gin.Context) {
in := &types.BaseIDReq{}
if err := x.Bind(c, in); err != nil {
response.Error(c, err)
return
}
cj, err := service.CronJobService.GetByID(ctxUtils.WrapCtx(c), in.ID, nil)
if err != nil {
response.Error(c, err)
return
}
data := types.CronJobObjDetailTool.NewFromBiz(cj)
response.Success(c, data)
}
func (x *cronJobController) Page(c *gin.Context) {
in := &types.CronJobPageReq{}
if err := x.Bind(c, in); err != nil {
response.Error(c, err)
return
}
cjs, total, err := service.CronJobService.Page(ctxUtils.WrapCtx(c), in, nil)
if err != nil {
response.Error(c, err)
return
}
data := types.CronJobObjDetailTool.NewSliceFromBizSlice(cjs)
response.SuccessWithPage(c, data, total)
}
func ExecJob(ctx context.Context, j *biz.CronJob) error {
if j == nil {
return errors.New("job is nil")
}
_, ok := Handlers.Get(j.JobCode)
if !ok {
return errors.New("handler not found")
}
return redSyncUtils.Sync(ctx, j.JobCode, func() {
if err := execTask(ctx, j); err != nil {
logger.Error(fmt.Sprintf("[cronJob][%s] exec err: %s", j.JobName, err.Error()), ctxUtils.CtxTraceIDField(ctx))
}
},
redsync.WithExpiry(5*time.Minute),
redsync.WithTries(10),
redsync.WithRetryDelayFunc(func(tries int) time.Duration {
return time.Duration(100+tries*20) * time.Millisecond
}),
redsync.WithDriftFactor(0.01),
redsync.WithTimeoutFactor(0.05))
}
func execTask(ctx context.Context, j *biz.CronJob) error {
if j == nil {
return errors.New("job not found")
}
fn, ok := Handlers.Get(j.JobCode)
if !ok {
return errors.New("handler not found")
}
startTime := time.Now()
job, err := service.CronJobService.GetNotRunning(ctx, j.ID)
if err != nil {
return err
}
if job == nil {
return nil
}
logger.Info(fmt.Sprintf("[cronJob][%s] start", j.JobName), ctxUtils.CtxTraceIDField(ctx))
cronJobRecord := &models.CronJobRecord{JobID: job.ID, BeginTime: xtime.Now()}
_ = dao.CronJobRecordDao.CreateSilent(ctx, cronJobRecord)
_ = dao.CronJobLogDao.CreateSilent(ctx, &models.CronJobLog{JobID: job.ID, Info: "job started."})
defer func() {
logger.Info(fmt.Sprintf("[cronJob][%s] end", j.JobName), logger.Any("cost", time.Since(startTime).String()), ctxUtils.CtxTraceIDField(ctx))
_ = dao.CronJobDao.UpdateByIDSilent(ctx, job.ID, &models.CronJob{State: enums.CronJob_State_NOT_RUNNING})
}()
_ = dao.CronJobDao.UpdateByIDSilent(ctx, job.ID, &models.CronJob{State: enums.CronJob_State_RUNNING})
if err = fn(ctx, j); err != nil {
_ = dao.CronJobLogDao.CreateSilent(ctx, &models.CronJobLog{JobID: job.ID, Info: "fail. err:" + err.Error()})
_ = dao.CronJobRecordDao.UpdateByIDSilent(ctx, cronJobRecord.ID, &models.CronJobRecord{EndTime: xtime.Now()})
return err
}
_ = dao.CronJobLogDao.CreateSilent(ctx, &models.CronJobLog{JobID: job.ID, Info: "success."})
_ = dao.CronJobRecordDao.UpdateByIDSilent(ctx, cronJobRecord.ID, &models.CronJobRecord{EndTime: xtime.Now()})
return nil
}
type Fn func(ctx context.Context, job *biz.CronJob) error
type handlers struct {
M map[string]Fn
}
var Handlers = &handlers{
M: make(map[string]Fn),
}
func (x *handlers) Get(key string) (Fn, bool) {
v, ok := x.M[key]
return v, ok
}
func (x *handlers) Add(key string, h Fn) {
x.M[key] = h
}
...@@ -2,119 +2,26 @@ package xcron ...@@ -2,119 +2,26 @@ package xcron
import ( import (
"context" "context"
"errors"
"fmt" "fmt"
"time"
"github.com/gin-gonic/gin"
"github.com/go-redsync/redsync/v4"
"gitlab.wanzhuangkj.com/tush/xpkg/config" "gitlab.wanzhuangkj.com/tush/xpkg/config"
"gitlab.wanzhuangkj.com/tush/xpkg/gin/response"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/eventbus" "gitlab.wanzhuangkj.com/tush/xpkg/pkg/eventbus"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/gocron" "gitlab.wanzhuangkj.com/tush/xpkg/pkg/gocron"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/logger" "gitlab.wanzhuangkj.com/tush/xpkg/pkg/logger"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/ctxUtils" "gitlab.wanzhuangkj.com/tush/xpkg/utils/ctxUtils"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/redSyncUtils"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/sliceUtils" "gitlab.wanzhuangkj.com/tush/xpkg/utils/sliceUtils"
"gitlab.wanzhuangkj.com/tush/xpkg/utils/xtime"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/biz" "gitlab.wanzhuangkj.com/tush/xpkg/xcron/biz"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/dao" "gitlab.wanzhuangkj.com/tush/xpkg/xcron/controller"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/enums" "gitlab.wanzhuangkj.com/tush/xpkg/xcron/enums"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/models"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/service" "gitlab.wanzhuangkj.com/tush/xpkg/xcron/service"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/types" "gitlab.wanzhuangkj.com/tush/xpkg/xcron/types"
) )
const ( func AddHandler(k string, h controller.Fn) {
PWD = "1937396767217688576" controller.Handlers.Add(k, h)
)
type cronJobController struct{}
var Instance cronJobController
func (x *cronJobController) StartByCode(c *gin.Context) {
req := &types.CronJobCodeReq{}
if err := c.ShouldBind(req); err != nil {
response.Error(c, err)
return
}
if req.Pwd != PWD {
response.Error(c, errors.New("pwd err"))
return
}
ctx := ctxUtils.WrapCtx(c)
job, err := service.CronJobService.GetByJobCode(ctx, req.JobCode, &biz.CronJobOpts{})
if err != nil {
response.Error(c, err)
return
}
if err := service.CronJobService.StartByID(ctx, job.ID, nil); err != nil {
response.Error(c, err)
return
}
gocron.Add(&gocron.Task{
TimeSpec: job.TimeSpec,
Name: job.JobCode,
Fn: func() {
ExecJob(context.TODO(), job)
},
})
response.Success(c)
}
func (x *cronJobController) StopByCode(c *gin.Context) {
req := &types.CronJobCodeReq{}
if err := c.ShouldBind(req); err != nil {
response.Error(c, err)
return
}
if req.Pwd != PWD {
response.Error(c, errors.New("pwd err"))
return
}
ctx := ctxUtils.WrapCtx(c)
job, err := service.CronJobService.GetByJobCode(ctx, req.JobCode, &biz.CronJobOpts{})
if err != nil {
response.Error(c, err)
return
}
if err := service.CronJobService.StopByID(ctx, job.ID, nil); err != nil {
response.Error(c, err)
return
}
gocron.DeleteTask(job.JobCode)
response.Success(c)
}
func (x *cronJobController) CallOnceByCode(c *gin.Context) {
req := &types.CronJobCodeReq{}
if err := c.ShouldBind(req); err != nil {
response.Error(c, err)
return
}
if req.Pwd != PWD {
response.Error(c, errors.New("pwd err"))
return
}
job, err := service.CronJobService.GetByJobCode(ctxUtils.WrapCtx(c), req.JobCode, &biz.CronJobOpts{})
if err != nil {
response.Error(c, err)
return
}
err = ExecJob(ctxUtils.WrapCtx(c), job)
if err != nil {
response.Error(c, err)
return
}
response.Success(c)
}
func AddHandler(k string, h Fn) {
Handlers.Add(k, h)
} }
func AddHandlerMap(handlers map[string]Fn) { func AddHandlerMap(handlers map[string]controller.Fn) {
for k, h := range handlers { for k, h := range handlers {
AddHandler(k, h) AddHandler(k, h)
} }
...@@ -138,14 +45,14 @@ func initCron() error { ...@@ -138,14 +45,14 @@ func initCron() error {
var tasks []*gocron.Task var tasks []*gocron.Task
for _, job := range cronJobs { for _, job := range cronJobs {
job.ParseJobInfo() job.ParseJobInfo()
_, ok := Handlers.Get(job.JobCode) _, ok := controller.Handlers.Get(job.JobCode)
if ok { if ok {
tsk := gocron.Task{ tsk := gocron.Task{
TimeSpec: job.TimeSpec, TimeSpec: job.TimeSpec,
Name: job.JobCode, Name: job.JobCode,
Fn: func() { Fn: func() {
ctx := context.WithValue(ctx, ctxUtils.HeaderXRequestIDKey, ctxUtils.GenerateTid()) ctx := context.WithValue(ctx, ctxUtils.HeaderXRequestIDKey, ctxUtils.GenerateTid())
if err := ExecJob(ctx, job); err != nil { if err := controller.ExecJob(ctx, job); err != nil {
logger.Error(fmt.Sprintf("[cron]exec err: %s", err.Error()), logger.Any("name", job.JobName), logger.Any("timeSpec", job.TimeSpec), logger.Any("args", job.Info), ctxUtils.CtxTraceIDField(ctx)) logger.Error(fmt.Sprintf("[cron]exec err: %s", err.Error()), logger.Any("name", job.JobName), logger.Any("timeSpec", job.TimeSpec), logger.Any("args", job.Info), ctxUtils.CtxTraceIDField(ctx))
} }
}, },
...@@ -167,81 +74,3 @@ func initCron() error { ...@@ -167,81 +74,3 @@ func initCron() error {
logger.Info("[cron] initialized", ctxUtils.CtxTraceIDField(ctx)) logger.Info("[cron] initialized", ctxUtils.CtxTraceIDField(ctx))
return nil return nil
} }
type Fn func(ctx context.Context, job *biz.CronJob) error
type handlers struct {
M map[string]Fn
}
var Handlers = &handlers{
M: make(map[string]Fn),
}
func (x *handlers) Get(key string) (Fn, bool) {
v, ok := x.M[key]
return v, ok
}
func (x *handlers) Add(key string, h Fn) {
x.M[key] = h
}
func ExecJob(ctx context.Context, j *biz.CronJob) error {
if j == nil {
return errors.New("job is nil")
}
_, ok := Handlers.Get(j.JobCode)
if !ok {
return errors.New("handler not found")
}
return redSyncUtils.Sync(ctx, j.JobCode, func() {
if err := execTask(ctx, j); err != nil {
logger.Error(fmt.Sprintf("[cronJob][%s] exec err: %s", j.JobName, err.Error()), ctxUtils.CtxTraceIDField(ctx))
}
},
redsync.WithExpiry(5*time.Minute),
redsync.WithTries(10),
redsync.WithRetryDelayFunc(func(tries int) time.Duration {
return time.Duration(100+tries*20) * time.Millisecond
}),
redsync.WithDriftFactor(0.01),
redsync.WithTimeoutFactor(0.05))
}
func execTask(ctx context.Context, j *biz.CronJob) error {
if j == nil {
return errors.New("job not found")
}
fn, ok := Handlers.Get(j.JobCode)
if !ok {
return errors.New("handler not found")
}
startTime := time.Now()
job, err := service.CronJobService.GetNotRunning(ctx, j.ID)
if err != nil {
return err
}
if job == nil {
return nil
}
logger.Info(fmt.Sprintf("[cronJob][%s] start", j.JobName), ctxUtils.CtxTraceIDField(ctx))
cronJobRecord := &models.CronJobRecord{JobID: job.ID, BeginTime: xtime.Now()}
_ = dao.CronJobRecordDao.CreateSilent(ctx, cronJobRecord)
_ = dao.CronJobLogDao.CreateSilent(ctx, &models.CronJobLog{JobID: job.ID, Info: "job started."})
defer func() {
logger.Info(fmt.Sprintf("[cronJob][%s] end", j.JobName), logger.Any("cost", time.Since(startTime).String()), ctxUtils.CtxTraceIDField(ctx))
_ = dao.CronJobDao.UpdateByIDSilent(ctx, job.ID, &models.CronJob{State: enums.CronJob_State_NOT_RUNNING})
}()
_ = dao.CronJobDao.UpdateByIDSilent(ctx, job.ID, &models.CronJob{State: enums.CronJob_State_RUNNING})
if err = fn(ctx, j); err != nil {
_ = dao.CronJobLogDao.CreateSilent(ctx, &models.CronJobLog{JobID: job.ID, Info: "fail. err:" + err.Error()})
_ = dao.CronJobRecordDao.UpdateByIDSilent(ctx, cronJobRecord.ID, &models.CronJobRecord{EndTime: xtime.Now()})
return err
}
_ = dao.CronJobLogDao.CreateSilent(ctx, &models.CronJobLog{JobID: job.ID, Info: "success."})
_ = dao.CronJobRecordDao.UpdateByIDSilent(ctx, cronJobRecord.ID, &models.CronJobRecord{EndTime: xtime.Now()})
return nil
}
...@@ -11,27 +11,16 @@ import ( ...@@ -11,27 +11,16 @@ import (
) )
func init() { func init() {
_ = eventbus.Eb.Subscribe(eventbus.TopicParseConfFinish, func(ctx context.Context) { _ = eventbus.Eb.Subscribe(eventbus.TopicCoreInitFinish, func(ctx context.Context) {
cronCfg := config.GetCron() cronCfg := config.GetCron()
if !cronCfg.Enable { if !cronCfg.Enable {
return return
} }
schema := "" schema := cronCfg.Schema
config.Read(func(c *config.Config) {
schema = c.Cron.Schema
})
dao.CronJobDao.ODao.XDB = &database.XDB{Schema: schema} dao.CronJobDao.ODao.XDB = &database.XDB{Schema: schema}
dao.CronJobLogDao.ODao.XDB = &database.XDB{Schema: schema} dao.CronJobLogDao.ODao.XDB = &database.XDB{Schema: schema}
dao.CronJobRecordDao.ODao.XDB = &database.XDB{Schema: schema} dao.CronJobRecordDao.ODao.XDB = &database.XDB{Schema: schema}
})
_ = eventbus.Eb.Subscribe(eventbus.TopicCacheInitFinish, func(ctx context.Context) {
cronCfg := config.GetCron()
if !cronCfg.Enable {
return
}
dao.CronJobDao.ODao.Cache = cache.NewCronJobCache(database.GetCacheType()) dao.CronJobDao.ODao.Cache = cache.NewCronJobCache(database.GetCacheType())
})
_ = eventbus.Eb.Subscribe(eventbus.TopicCoreInitFinish, func(ctx context.Context) {
initCron() initCron()
}) })
} }
package xcron
import (
"github.com/gin-gonic/gin"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/controller"
)
var CronRouterApis []func(r *gin.RouterGroup) //
func init() {
CronRouterApis = append(CronRouterApis, func(group *gin.RouterGroup) {
cronJobRouter(group)
})
}
func cronJobRouter(group *gin.RouterGroup) {
c := controller.CronJobController
g := group.Group("/cronJob")
g.POST("/startByCode", c.StartByCode)
g.POST("/stopByCode", c.StopByCode)
g.POST("/callByCode", c.CallOnceByCode)
g.POST("/create", c.Create)
g.DELETE("/deleteByID", c.DeleteByID)
g.DELETE("/updateByID", c.UpdateByID)
g.POST("/getByID", c.GetByID)
g.POST("/page", c.Page)
}
...@@ -99,3 +99,22 @@ type CronJobCodeReq struct { ...@@ -99,3 +99,22 @@ type CronJobCodeReq struct {
JobCode string `json:"jobCode" form:"jobCode" validate:"required" example:"autoConfirmOrdersTask"` // 定时任务code JobCode string `json:"jobCode" form:"jobCode" validate:"required" example:"autoConfirmOrdersTask"` // 定时任务code
Pwd string `json:"pwd" form:"pwd" validate:"required" example:"abc123"` // Pwd string `json:"pwd" form:"pwd" validate:"required" example:"abc123"` //
} }
var CronJobObjDetailTool = CronJobObjDetail{}
func (x CronJobObjDetail) NewFromBiz(cj *biz.CronJob) *CronJobObjDetail {
if cj == nil {
return nil
}
return &CronJobObjDetail{
CronJob: *cj,
}
}
func (x CronJobObjDetail) NewSliceFromBizSlice(cjs []*biz.CronJob) []*CronJobObjDetail {
ret := make([]*CronJobObjDetail, 0, len(cjs))
for i := range cjs {
ret = append(ret, x.NewFromBiz(cjs[i]))
}
return ret
}
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论