提交 f591f912 authored 作者: mooncake9527's avatar mooncake9527

update

上级 6bc8686b
...@@ -72,6 +72,14 @@ func GetOssConfig() *oss.AliOssConfig { ...@@ -72,6 +72,14 @@ func GetOssConfig() *oss.AliOssConfig {
return ossCfg return ossCfg
} }
func GetCron() *Cron {
cron := Cron{}
Read(func(c *Config) {
copier.Copy(&cron, &c.Cron)
})
return &cron
}
func Read(f func(c *Config)) { func Read(f func(c *Config)) {
if Cfg == nil { if Cfg == nil {
panic("config is nil, please call config.Init() first") panic("config is nil, please call config.Init() first")
......
...@@ -110,14 +110,22 @@ func (x *cronJobController) CallOnceByCode(c *gin.Context) { ...@@ -110,14 +110,22 @@ func (x *cronJobController) CallOnceByCode(c *gin.Context) {
response.Success(c) response.Success(c)
} }
func Init(schema string, handlers map[string]Fn) error { func AddHandler(k string, h Fn) {
dao.Init(schema) Handlers.Add(k, h)
}
ctx := context.TODO() func AddHandlerMap(handlers map[string]Fn) {
TaskRegisterMaps = handlers for k, h := range handlers {
if !config.CronOpen() { AddHandler(k, h)
}
}
func initCron() error {
cronCfg := config.GetCron()
if !cronCfg.Enable {
return nil return nil
} }
ctx := context.TODO()
cronJobs, err := service.CronJobService.GetSliceByEnable(ctx, enums.CronJob_Enable_OPEN, &biz.CronJobOpts{}) cronJobs, err := service.CronJobService.GetSliceByEnable(ctx, enums.CronJob_Enable_OPEN, &biz.CronJobOpts{})
if err != nil { if err != nil {
return err return err
...@@ -130,7 +138,7 @@ func Init(schema string, handlers map[string]Fn) error { ...@@ -130,7 +138,7 @@ func Init(schema string, handlers map[string]Fn) error {
var tasks []*gocron.Task var tasks []*gocron.Task
for _, job := range cronJobs { for _, job := range cronJobs {
job.ParseJobInfo() job.ParseJobInfo()
_, ok := handlers[job.JobCode] _, ok := Handlers.Get(job.JobCode)
if ok { if ok {
tsk := gocron.Task{ tsk := gocron.Task{
TimeSpec: job.TimeSpec, TimeSpec: job.TimeSpec,
...@@ -162,13 +170,28 @@ func Init(schema string, handlers map[string]Fn) error { ...@@ -162,13 +170,28 @@ func Init(schema string, handlers map[string]Fn) error {
type Fn func(ctx context.Context, job *biz.CronJob) error type Fn func(ctx context.Context, job *biz.CronJob) error
var TaskRegisterMaps = map[string]Fn{} 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 { func ExecJob(ctx context.Context, j *biz.CronJob) error {
if j == nil { if j == nil {
return errors.New("job is nil") return errors.New("job is nil")
} }
_, ok := TaskRegisterMaps[j.JobCode] _, ok := Handlers.Get(j.JobCode)
if !ok { if !ok {
return errors.New("handler not found") return errors.New("handler not found")
} }
...@@ -190,7 +213,7 @@ func execTask(ctx context.Context, j *biz.CronJob) error { ...@@ -190,7 +213,7 @@ func execTask(ctx context.Context, j *biz.CronJob) error {
if j == nil { if j == nil {
return errors.New("job not found") return errors.New("job not found")
} }
fn, ok := TaskRegisterMaps[j.JobCode] fn, ok := Handlers.Get(j.JobCode)
if !ok { if !ok {
return errors.New("handler not found") return errors.New("handler not found")
} }
......
package dao
import (
"context"
"gitlab.wanzhuangkj.com/tush/xpkg/database"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/eventbus"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/cache"
)
var (
Schema string
)
func Init(schema string) {
Schema = schema
_ = eventbus.Eb.Subscribe(eventbus.TopicCacheInitFinish, func(ctx context.Context) {
CronJobDao.ODao.Cache = cache.NewCronJobCache(database.GetCacheType())
})
}
package dao package xcron
import ( import (
"context" "context"
...@@ -6,16 +6,32 @@ import ( ...@@ -6,16 +6,32 @@ import (
"gitlab.wanzhuangkj.com/tush/xpkg/config" "gitlab.wanzhuangkj.com/tush/xpkg/config"
"gitlab.wanzhuangkj.com/tush/xpkg/database" "gitlab.wanzhuangkj.com/tush/xpkg/database"
"gitlab.wanzhuangkj.com/tush/xpkg/pkg/eventbus" "gitlab.wanzhuangkj.com/tush/xpkg/pkg/eventbus"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/cache"
"gitlab.wanzhuangkj.com/tush/xpkg/xcron/dao"
) )
func init() { func init() {
_ = eventbus.Eb.Subscribe(eventbus.TopicParseConfFinish, func(ctx context.Context) { _ = eventbus.Eb.Subscribe(eventbus.TopicParseConfFinish, func(ctx context.Context) {
cronCfg := config.GetCron()
if !cronCfg.Enable {
return
}
schema := "" schema := ""
config.Read(func(c *config.Config) { config.Read(func(c *config.Config) {
schema = c.Cron.Schema schema = c.Cron.Schema
}) })
CronJobDao.ODao.XDB = &database.XDB{Schema: schema} dao.CronJobDao.ODao.XDB = &database.XDB{Schema: schema}
CronJobLogDao.ODao.XDB = &database.XDB{Schema: schema} dao.CronJobLogDao.ODao.XDB = &database.XDB{Schema: schema}
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())
})
_ = eventbus.Eb.Subscribe(eventbus.TopicCoreInitFinish, func(ctx context.Context) {
initCron()
}) })
} }
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论