mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-04 15:06:37 +08:00
refactor(core): decouple gin from pkg/util and reduce code duplication
This commit is contained in:
@@ -139,34 +139,7 @@ func (p *Plugin) Start(_ context.Context) error {
|
||||
}
|
||||
|
||||
if p.scheduler == nil {
|
||||
opts := p.schedulerOpts
|
||||
if opts == nil {
|
||||
opts = &asynq.SchedulerOpts{
|
||||
Location: p.location,
|
||||
}
|
||||
} else if opts.Location == nil && p.location != nil {
|
||||
opts.Location = p.location
|
||||
}
|
||||
|
||||
opt := p.redisOpt
|
||||
if opt == nil {
|
||||
if RedisOpt != nil {
|
||||
opt = RedisOpt
|
||||
} else {
|
||||
redisCfg := config.Config.Redis
|
||||
addr := "127.0.0.1:6379"
|
||||
if len(redisCfg.Addrs) > 0 && redisCfg.Addrs[0] != "" {
|
||||
addr = redisCfg.Addrs[0]
|
||||
}
|
||||
opt = asynq.RedisClientOpt{
|
||||
Addr: addr,
|
||||
Username: redisCfg.Username,
|
||||
Password: redisCfg.Password,
|
||||
DB: redisCfg.DB,
|
||||
}
|
||||
}
|
||||
}
|
||||
p.scheduler = asynq.NewScheduler(opt, opts)
|
||||
p.scheduler = p.initScheduler()
|
||||
}
|
||||
|
||||
if p.coreCtx != nil && p.coreCtx.Schedules() != nil {
|
||||
@@ -268,3 +241,37 @@ func buildAsynqOptions(opts map[string]any) []asynq.Option {
|
||||
|
||||
return res
|
||||
}
|
||||
|
||||
func (p *Plugin) initScheduler() *asynq.Scheduler {
|
||||
opts := p.schedulerOpts
|
||||
if opts == nil {
|
||||
opts = &asynq.SchedulerOpts{
|
||||
Location: p.location,
|
||||
}
|
||||
} else if opts.Location == nil && p.location != nil {
|
||||
opts.Location = p.location
|
||||
}
|
||||
|
||||
opt := p.resolveRedisOpt()
|
||||
return asynq.NewScheduler(opt, opts)
|
||||
}
|
||||
|
||||
func (p *Plugin) resolveRedisOpt() asynq.RedisConnOpt {
|
||||
if p.redisOpt != nil {
|
||||
return p.redisOpt
|
||||
}
|
||||
if RedisOpt != nil {
|
||||
return RedisOpt
|
||||
}
|
||||
redisCfg := config.Config.Redis
|
||||
addr := "127.0.0.1:6379"
|
||||
if len(redisCfg.Addrs) > 0 && redisCfg.Addrs[0] != "" {
|
||||
addr = redisCfg.Addrs[0]
|
||||
}
|
||||
return asynq.RedisClientOpt{
|
||||
Addr: addr,
|
||||
Username: redisCfg.Username,
|
||||
Password: redisCfg.Password,
|
||||
DB: redisCfg.DB,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -111,21 +111,7 @@ func ReloadScheduler() error {
|
||||
// 4. 遍历并注册任务
|
||||
taskSvc := getTaskService()
|
||||
for _, s := range schedules {
|
||||
taskName := s.TaskType
|
||||
maxRetry := 3
|
||||
queue := "default"
|
||||
|
||||
if taskSvc != nil {
|
||||
if meta, ok := taskSvc.GetTaskMeta(s.TaskType); ok {
|
||||
taskName = meta.Name
|
||||
if meta.MaxRetry > 0 {
|
||||
maxRetry = meta.MaxRetry
|
||||
}
|
||||
if meta.Queue != "" {
|
||||
queue = meta.Queue
|
||||
}
|
||||
}
|
||||
}
|
||||
taskName, maxRetry, queue := resolveTaskScheduleMeta(taskSvc, s.TaskType)
|
||||
|
||||
// 构造 Asynq 载荷
|
||||
t := asynq.NewTask(taskName, []byte(s.Payload))
|
||||
@@ -160,3 +146,24 @@ func waitForStop(done, signals <-chan struct{}) bool {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
func resolveTaskScheduleMeta(taskSvc contracts.TaskService, taskType string) (name string, maxRetry int, queue string) {
|
||||
name = taskType
|
||||
maxRetry = 3
|
||||
queue = "default"
|
||||
if taskSvc == nil {
|
||||
return
|
||||
}
|
||||
meta, ok := taskSvc.GetTaskMeta(taskType)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
name = meta.Name
|
||||
if meta.MaxRetry > 0 {
|
||||
maxRetry = meta.MaxRetry
|
||||
}
|
||||
if meta.Queue != "" {
|
||||
queue = meta.Queue
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
@@ -366,27 +366,8 @@ func (s *taskServiceImpl) ListExecutions(ctx context.Context, taskType, status s
|
||||
return nil, 0, err
|
||||
}
|
||||
res := make([]contracts.TaskExecutionDTO, 0, len(rows))
|
||||
for _, r := range rows {
|
||||
res = append(res, contracts.TaskExecutionDTO{
|
||||
ID: r.ID,
|
||||
TaskID: r.TaskID,
|
||||
TaskType: r.TaskType,
|
||||
TaskName: r.TaskName,
|
||||
Status: string(r.Status),
|
||||
Retryable: r.Retryable,
|
||||
MaxRetry: r.MaxRetry,
|
||||
RetryCount: r.RetryCount,
|
||||
Log: r.Log,
|
||||
ErrorMessage: r.ErrorMessage,
|
||||
Result: r.Result,
|
||||
StartedAt: r.StartedAt,
|
||||
FinishedAt: r.FinishedAt,
|
||||
Duration: r.Duration,
|
||||
Payload: r.Payload,
|
||||
TriggeredBy: r.TriggeredBy,
|
||||
CreatedAt: r.CreatedAt,
|
||||
UpdatedAt: r.UpdatedAt,
|
||||
})
|
||||
for i := range rows {
|
||||
res = append(res, toTaskExecutionDTO(&rows[i]))
|
||||
}
|
||||
return res, total, nil
|
||||
}
|
||||
@@ -416,7 +397,12 @@ func (s *taskServiceImpl) GetExecution(ctx context.Context, id uint64) (*contrac
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &contracts.TaskExecutionDTO{
|
||||
dto := toTaskExecutionDTO(exec)
|
||||
return &dto, nil
|
||||
}
|
||||
|
||||
func toTaskExecutionDTO(exec *TaskExecution) contracts.TaskExecutionDTO {
|
||||
return contracts.TaskExecutionDTO{
|
||||
ID: exec.ID,
|
||||
TaskID: exec.TaskID,
|
||||
TaskType: exec.TaskType,
|
||||
@@ -435,5 +421,5 @@ func (s *taskServiceImpl) GetExecution(ctx context.Context, id uint64) (*contrac
|
||||
TriggeredBy: exec.TriggeredBy,
|
||||
CreatedAt: exec.CreatedAt,
|
||||
UpdatedAt: exec.UpdatedAt,
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -78,7 +79,7 @@ func (s *inprocScheduler) registerJob(ctx context.Context, def extpoints.Schedul
|
||||
spec := def.Spec
|
||||
taskType := def.TaskType
|
||||
|
||||
fields := len(cronFields(spec))
|
||||
fields := len(strings.Fields(spec))
|
||||
cronSpec := spec
|
||||
if fields == standardCronFields {
|
||||
cronSpec = "0 " + spec
|
||||
@@ -141,22 +142,3 @@ func invokeHandler(ctx context.Context, handler any, payload []byte) error {
|
||||
return fmt.Errorf("unsupported handler type: %T", handler)
|
||||
}
|
||||
}
|
||||
|
||||
func cronFields(s string) []string {
|
||||
var fields []string
|
||||
var current []rune
|
||||
for _, r := range s {
|
||||
if r == ' ' || r == '\t' {
|
||||
if len(current) > 0 {
|
||||
fields = append(fields, string(current))
|
||||
current = nil
|
||||
}
|
||||
} else {
|
||||
current = append(current, r)
|
||||
}
|
||||
}
|
||||
if len(current) > 0 {
|
||||
fields = append(fields, string(current))
|
||||
}
|
||||
return fields
|
||||
}
|
||||
|
||||
@@ -101,7 +101,7 @@ func (q *InprocQueue) Start(ctx context.Context) {
|
||||
q.wg.Add(1)
|
||||
util.Go(func() {
|
||||
defer q.wg.Done()
|
||||
q.workerLoop()
|
||||
q.workerLoop(ctx)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -128,21 +128,23 @@ func (q *InprocQueue) Stop(ctx context.Context) error {
|
||||
}
|
||||
}
|
||||
|
||||
func (q *InprocQueue) workerLoop() {
|
||||
func (q *InprocQueue) workerLoop(ctx context.Context) {
|
||||
for {
|
||||
select {
|
||||
case <-q.stopCh:
|
||||
return
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case msg, ok := <-q.queue:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
q.executeTask(msg)
|
||||
q.executeTask(ctx, msg)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (q *InprocQueue) executeTask(msg TaskMessage) {
|
||||
func (q *InprocQueue) executeTask(ctx context.Context, msg TaskMessage) {
|
||||
if q.taskReg == nil {
|
||||
return
|
||||
}
|
||||
@@ -157,7 +159,7 @@ func (q *InprocQueue) executeTask(msg TaskMessage) {
|
||||
timeout = 5 * time.Minute
|
||||
}
|
||||
|
||||
taskCtx, cancel := context.WithTimeout(q.baseCtx, timeout)
|
||||
taskCtx, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
|
||||
err := invokeHandler(taskCtx, td.Handler, msg.Payload)
|
||||
@@ -165,7 +167,13 @@ func (q *InprocQueue) executeTask(msg TaskMessage) {
|
||||
msg.RetryLeft--
|
||||
// Retry with backoff
|
||||
util.Go(func() {
|
||||
time.Sleep(defaultRetryBackoff)
|
||||
select {
|
||||
case <-time.After(defaultRetryBackoff):
|
||||
case <-q.stopCh:
|
||||
return
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
if q.running.Load() {
|
||||
select {
|
||||
case q.queue <- msg:
|
||||
|
||||
Reference in New Issue
Block a user