mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-11 01:36:37 +08:00
refactor(drivers): decouple inproc cron from inproc worker and satisfy linter
This commit is contained in:
@@ -9,6 +9,7 @@ import (
|
|||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"Wavelet/core"
|
"Wavelet/core"
|
||||||
|
"Wavelet/core/contracts"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Plugin implements core.Plugin and core.Driver for in-process cron job scheduling.
|
// Plugin implements core.Plugin and core.Driver for in-process cron job scheduling.
|
||||||
@@ -62,7 +63,8 @@ func (p *Plugin) Start(_ context.Context) error {
|
|||||||
defer p.mu.Unlock()
|
defer p.mu.Unlock()
|
||||||
|
|
||||||
if p.scheduler == nil {
|
if p.scheduler == nil {
|
||||||
p.scheduler = newInprocScheduler(p.coreCtx.Schedules(), p.coreCtx.Tasks())
|
taskSvc, _ := core.Inject[contracts.TaskService](p.coreCtx)
|
||||||
|
p.scheduler = newInprocScheduler(p.coreCtx.Schedules(), p.coreCtx.Tasks(), taskSvc)
|
||||||
}
|
}
|
||||||
|
|
||||||
return p.scheduler.Start()
|
return p.scheduler.Start()
|
||||||
|
|||||||
@@ -14,17 +14,11 @@ import (
|
|||||||
|
|
||||||
"Wavelet/core"
|
"Wavelet/core"
|
||||||
"Wavelet/plugins/drivers/driver_inproc_cron"
|
"Wavelet/plugins/drivers/driver_inproc_cron"
|
||||||
"Wavelet/plugins/drivers/driver_inproc_worker"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestInprocCronPlugin(t *testing.T) {
|
func TestInprocCronPlugin(t *testing.T) {
|
||||||
ctx := core.NewContext(context.Background())
|
ctx := core.NewContext(context.Background())
|
||||||
|
|
||||||
workerPlugin := driver_inproc_worker.New()
|
|
||||||
require.NoError(t, workerPlugin.Apply(ctx))
|
|
||||||
require.NoError(t, workerPlugin.Start(context.Background()))
|
|
||||||
defer func() { _ = workerPlugin.Stop(context.Background()) }()
|
|
||||||
|
|
||||||
cronPlugin := driver_inproc_cron.New()
|
cronPlugin := driver_inproc_cron.New()
|
||||||
assert.Equal(t, "driver_inproc_cron", cronPlugin.Name())
|
assert.Equal(t, "driver_inproc_cron", cronPlugin.Name())
|
||||||
assert.Equal(t, core.DriverTypeScheduler, cronPlugin.Type())
|
assert.Equal(t, core.DriverTypeScheduler, cronPlugin.Type())
|
||||||
|
|||||||
@@ -6,12 +6,17 @@ package driver_inproc_cron
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
"sync"
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/robfig/cron/v3"
|
||||||
|
|
||||||
|
"Wavelet/core/contracts"
|
||||||
"Wavelet/core/extpoints"
|
"Wavelet/core/extpoints"
|
||||||
"Wavelet/pkg/logger"
|
"Wavelet/pkg/logger"
|
||||||
"Wavelet/plugins/drivers/driver_inproc_worker"
|
"Wavelet/pkg/util"
|
||||||
"github.com/robfig/cron/v3"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type inprocScheduler struct {
|
type inprocScheduler struct {
|
||||||
@@ -19,14 +24,16 @@ type inprocScheduler struct {
|
|||||||
cronRunner *cron.Cron
|
cronRunner *cron.Cron
|
||||||
scheduleReg extpoints.ScheduleExtension
|
scheduleReg extpoints.ScheduleExtension
|
||||||
taskReg extpoints.TaskExtension
|
taskReg extpoints.TaskExtension
|
||||||
|
taskSvc contracts.TaskService
|
||||||
running bool
|
running bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func newInprocScheduler(scheduleReg extpoints.ScheduleExtension, taskReg extpoints.TaskExtension) *inprocScheduler {
|
func newInprocScheduler(scheduleReg extpoints.ScheduleExtension, taskReg extpoints.TaskExtension, taskSvc contracts.TaskService) *inprocScheduler {
|
||||||
return &inprocScheduler{
|
return &inprocScheduler{
|
||||||
cronRunner: cron.New(cron.WithSeconds()),
|
cronRunner: cron.New(cron.WithSeconds()),
|
||||||
scheduleReg: scheduleReg,
|
scheduleReg: scheduleReg,
|
||||||
taskReg: taskReg,
|
taskReg: taskReg,
|
||||||
|
taskSvc: taskSvc,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -62,14 +69,15 @@ func (s *inprocScheduler) Stop() {
|
|||||||
s.running = false
|
s.running = false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const standardCronFields = 5
|
||||||
|
|
||||||
func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) {
|
func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) {
|
||||||
spec := def.Spec
|
spec := def.Spec
|
||||||
taskType := def.TaskType
|
taskType := def.TaskType
|
||||||
|
|
||||||
// Support 5-field cron expression by prepending "0 " for 6-field seconds parser
|
|
||||||
fields := len(cronFields(spec))
|
fields := len(cronFields(spec))
|
||||||
cronSpec := spec
|
cronSpec := spec
|
||||||
if fields == 5 {
|
if fields == standardCronFields {
|
||||||
cronSpec = "0 " + spec
|
cronSpec = "0 " + spec
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -86,9 +94,25 @@ func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
_, err := s.cronRunner.AddFunc(cronSpec, func() {
|
_, err := s.cronRunner.AddFunc(cronSpec, func() {
|
||||||
_, dispatchErr := driver_inproc_worker.DispatchTask(context.Background(), taskType, payloadBytes, "inproc_cron")
|
if s.taskSvc != nil {
|
||||||
if dispatchErr != nil {
|
if _, dispatchErr := s.taskSvc.Dispatch(context.Background(), taskType, payloadBytes, "inproc_cron"); dispatchErr != nil {
|
||||||
logger.ErrorF(context.Background(), "driver_inproc_cron: dispatch task %q failed: %v", taskType, dispatchErr)
|
logger.ErrorF(context.Background(), "driver_inproc_cron: dispatch task %q failed: %v", taskType, dispatchErr)
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
if s.taskReg != nil {
|
||||||
|
if td, ok := s.taskReg.Get(taskType); ok {
|
||||||
|
util.Go(func() {
|
||||||
|
timeout := td.Timeout
|
||||||
|
if timeout <= 0 {
|
||||||
|
timeout = 5 * time.Minute
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), timeout)
|
||||||
|
defer cancel()
|
||||||
|
_ = invokeHandler(ctx, td.Handler, payloadBytes)
|
||||||
|
})
|
||||||
|
}
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -96,6 +120,25 @@ func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func invokeHandler(ctx context.Context, handler any, payload []byte) error {
|
||||||
|
if handler == nil {
|
||||||
|
return errors.New("nil task handler")
|
||||||
|
}
|
||||||
|
|
||||||
|
switch fn := handler.(type) {
|
||||||
|
case func(context.Context, []byte) error:
|
||||||
|
return fn(ctx, payload)
|
||||||
|
case func(context.Context) error:
|
||||||
|
return fn(ctx)
|
||||||
|
case func([]byte) error:
|
||||||
|
return fn(payload)
|
||||||
|
case func() error:
|
||||||
|
return fn()
|
||||||
|
default:
|
||||||
|
return fmt.Errorf("unsupported handler type: %T", handler)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func cronFields(s string) []string {
|
func cronFields(s string) []string {
|
||||||
var fields []string
|
var fields []string
|
||||||
var current []rune
|
var current []rune
|
||||||
|
|||||||
@@ -16,6 +16,8 @@ import (
|
|||||||
"Wavelet/pkg/util"
|
"Wavelet/pkg/util"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const defaultRetryBackoff = 500 * time.Millisecond
|
||||||
|
|
||||||
// TaskMessage represents an in-process task item in the queue.
|
// TaskMessage represents an in-process task item in the queue.
|
||||||
type TaskMessage struct {
|
type TaskMessage struct {
|
||||||
ID string
|
ID string
|
||||||
@@ -28,7 +30,6 @@ type TaskMessage struct {
|
|||||||
|
|
||||||
// InprocQueue manages in-memory task queuing and worker pool execution.
|
// InprocQueue manages in-memory task queuing and worker pool execution.
|
||||||
type InprocQueue struct {
|
type InprocQueue struct {
|
||||||
mu sync.RWMutex
|
|
||||||
concurrency int
|
concurrency int
|
||||||
queue chan TaskMessage
|
queue chan TaskMessage
|
||||||
taskReg extpoints.TaskExtension
|
taskReg extpoints.TaskExtension
|
||||||
@@ -157,7 +158,7 @@ func (q *InprocQueue) executeTask(msg TaskMessage) {
|
|||||||
msg.RetryLeft--
|
msg.RetryLeft--
|
||||||
// Retry with backoff
|
// Retry with backoff
|
||||||
util.Go(func() {
|
util.Go(func() {
|
||||||
time.Sleep(500 * time.Millisecond)
|
time.Sleep(defaultRetryBackoff)
|
||||||
if q.running.Load() {
|
if q.running.Load() {
|
||||||
select {
|
select {
|
||||||
case q.queue <- msg:
|
case q.queue <- msg:
|
||||||
|
|||||||
@@ -25,7 +25,7 @@ func (s *inprocTaskService) Dispatch(ctx context.Context, taskType string, paylo
|
|||||||
return DispatchTask(ctx, taskType, payload, triggeredBy)
|
return DispatchTask(ctx, taskType, payload, triggeredBy)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *inprocTaskService) Retry(ctx context.Context, id uint64) (string, error) {
|
func (s *inprocTaskService) Retry(_ context.Context, id uint64) (string, error) {
|
||||||
return fmt.Sprintf("inproc_retry_%d", id), nil
|
return fmt.Sprintf("inproc_retry_%d", id), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -37,7 +37,7 @@ func newMemoryCacheService(capacity int, events *core.EventBus) (*memoryCacheSer
|
|||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *memoryCacheService) Get(ctx context.Context, key string, target any) error {
|
func (s *memoryCacheService) Get(_ context.Context, key string, target any) error {
|
||||||
if entry, ok := s.ramCache.GetIfPresent(key); ok {
|
if entry, ok := s.ramCache.GetIfPresent(key); ok {
|
||||||
if entry.expireAt.IsZero() || time.Now().Before(entry.expireAt) {
|
if entry.expireAt.IsZero() || time.Now().Before(entry.expireAt) {
|
||||||
return json.Unmarshal(entry.data, target)
|
return json.Unmarshal(entry.data, target)
|
||||||
@@ -48,7 +48,7 @@ func (s *memoryCacheService) Get(ctx context.Context, key string, target any) er
|
|||||||
return contracts.ErrCacheMiss
|
return contracts.ErrCacheMiss
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *memoryCacheService) Set(ctx context.Context, key string, value any, ttl time.Duration) error {
|
func (s *memoryCacheService) Set(_ context.Context, key string, value any, ttl time.Duration) error {
|
||||||
data, err := json.Marshal(value)
|
data, err := json.Marshal(value)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
|
|||||||
Reference in New Issue
Block a user