diff --git a/backend/plugins/drivers/driver_inproc_cron/plugin.go b/backend/plugins/drivers/driver_inproc_cron/plugin.go index 19217521..2279b219 100644 --- a/backend/plugins/drivers/driver_inproc_cron/plugin.go +++ b/backend/plugins/drivers/driver_inproc_cron/plugin.go @@ -9,6 +9,7 @@ import ( "sync" "Wavelet/core" + "Wavelet/core/contracts" ) // 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() 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() diff --git a/backend/plugins/drivers/driver_inproc_cron/plugin_test.go b/backend/plugins/drivers/driver_inproc_cron/plugin_test.go index 31e75e62..09137b71 100644 --- a/backend/plugins/drivers/driver_inproc_cron/plugin_test.go +++ b/backend/plugins/drivers/driver_inproc_cron/plugin_test.go @@ -14,17 +14,11 @@ import ( "Wavelet/core" "Wavelet/plugins/drivers/driver_inproc_cron" - "Wavelet/plugins/drivers/driver_inproc_worker" ) func TestInprocCronPlugin(t *testing.T) { 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() assert.Equal(t, "driver_inproc_cron", cronPlugin.Name()) assert.Equal(t, core.DriverTypeScheduler, cronPlugin.Type()) diff --git a/backend/plugins/drivers/driver_inproc_cron/scheduler.go b/backend/plugins/drivers/driver_inproc_cron/scheduler.go index 6a0fc509..a8e23966 100644 --- a/backend/plugins/drivers/driver_inproc_cron/scheduler.go +++ b/backend/plugins/drivers/driver_inproc_cron/scheduler.go @@ -6,12 +6,17 @@ package driver_inproc_cron import ( "context" "encoding/json" + "errors" + "fmt" "sync" + "time" + "github.com/robfig/cron/v3" + + "Wavelet/core/contracts" "Wavelet/core/extpoints" "Wavelet/pkg/logger" - "Wavelet/plugins/drivers/driver_inproc_worker" - "github.com/robfig/cron/v3" + "Wavelet/pkg/util" ) type inprocScheduler struct { @@ -19,14 +24,16 @@ type inprocScheduler struct { cronRunner *cron.Cron scheduleReg extpoints.ScheduleExtension taskReg extpoints.TaskExtension + taskSvc contracts.TaskService 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{ cronRunner: cron.New(cron.WithSeconds()), scheduleReg: scheduleReg, taskReg: taskReg, + taskSvc: taskSvc, } } @@ -62,14 +69,15 @@ func (s *inprocScheduler) Stop() { s.running = false } +const standardCronFields = 5 + func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) { spec := def.Spec taskType := def.TaskType - // Support 5-field cron expression by prepending "0 " for 6-field seconds parser fields := len(cronFields(spec)) cronSpec := spec - if fields == 5 { + if fields == standardCronFields { cronSpec = "0 " + spec } @@ -86,9 +94,25 @@ func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) { } _, err := s.cronRunner.AddFunc(cronSpec, func() { - _, dispatchErr := driver_inproc_worker.DispatchTask(context.Background(), taskType, payloadBytes, "inproc_cron") - if dispatchErr != nil { - logger.ErrorF(context.Background(), "driver_inproc_cron: dispatch task %q failed: %v", taskType, dispatchErr) + if s.taskSvc != 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) + } + 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 { @@ -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 { var fields []string var current []rune diff --git a/backend/plugins/drivers/driver_inproc_worker/executor.go b/backend/plugins/drivers/driver_inproc_worker/executor.go index cff4fe12..6b3cb5d9 100644 --- a/backend/plugins/drivers/driver_inproc_worker/executor.go +++ b/backend/plugins/drivers/driver_inproc_worker/executor.go @@ -16,6 +16,8 @@ import ( "Wavelet/pkg/util" ) +const defaultRetryBackoff = 500 * time.Millisecond + // TaskMessage represents an in-process task item in the queue. type TaskMessage struct { ID string @@ -28,7 +30,6 @@ type TaskMessage struct { // InprocQueue manages in-memory task queuing and worker pool execution. type InprocQueue struct { - mu sync.RWMutex concurrency int queue chan TaskMessage taskReg extpoints.TaskExtension @@ -157,7 +158,7 @@ func (q *InprocQueue) executeTask(msg TaskMessage) { msg.RetryLeft-- // Retry with backoff util.Go(func() { - time.Sleep(500 * time.Millisecond) + time.Sleep(defaultRetryBackoff) if q.running.Load() { select { case q.queue <- msg: diff --git a/backend/plugins/drivers/driver_inproc_worker/task_service.go b/backend/plugins/drivers/driver_inproc_worker/task_service.go index 1721393a..9726bf3c 100644 --- a/backend/plugins/drivers/driver_inproc_worker/task_service.go +++ b/backend/plugins/drivers/driver_inproc_worker/task_service.go @@ -25,7 +25,7 @@ func (s *inprocTaskService) Dispatch(ctx context.Context, taskType string, paylo 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 } diff --git a/backend/plugins/infra/cache_memory/cache.go b/backend/plugins/infra/cache_memory/cache.go index b617541f..9639975b 100644 --- a/backend/plugins/infra/cache_memory/cache.go +++ b/backend/plugins/infra/cache_memory/cache.go @@ -37,7 +37,7 @@ func newMemoryCacheService(capacity int, events *core.EventBus) (*memoryCacheSer }, 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.expireAt.IsZero() || time.Now().Before(entry.expireAt) { 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 } -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) if err != nil { return err