diff --git a/backend/plugins/drivers/driver_inproc_cron/plugin.go b/backend/plugins/drivers/driver_inproc_cron/plugin.go new file mode 100644 index 00000000..19217521 --- /dev/null +++ b/backend/plugins/drivers/driver_inproc_cron/plugin.go @@ -0,0 +1,81 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +// Package driver_inproc_cron provides the zero-dependency in-process cron scheduler driver plugin for Cordis. +package driver_inproc_cron + +import ( + "context" + "sync" + + "Wavelet/core" +) + +// Plugin implements core.Plugin and core.Driver for in-process cron job scheduling. +type Plugin struct { + mu sync.RWMutex + coreCtx *core.Context + scheduler *inprocScheduler +} + +// New creates a new in-process cron scheduler driver plugin. +func New() *Plugin { + return &Plugin{} +} + +// Name returns the unique identifier of the plugin. +func (p *Plugin) Name() string { + return "driver_inproc_cron" +} + +// Manifest returns the plugin metadata. +func (p *Plugin) Manifest() core.Manifest { + return core.Manifest{ + Name: "driver_inproc_cron", + Version: "1.0.0", + Description: "Zero-dependency in-process cron scheduler driver plugin", + Author: "Wavelet Team", + } +} + +// Apply registers the scheduler driver into the Context. +func (p *Plugin) Apply(ctx *core.Context) error { + p.mu.Lock() + p.coreCtx = ctx + p.mu.Unlock() + + ctx.OnDispose(func() error { + return p.Stop(context.Background()) + }) + + return ctx.RegisterDriver(p) +} + +// Type returns DriverTypeScheduler. +func (p *Plugin) Type() core.DriverType { + return core.DriverTypeScheduler +} + +// Start boots the in-process cron scheduler. +func (p *Plugin) Start(_ context.Context) error { + p.mu.Lock() + defer p.mu.Unlock() + + if p.scheduler == nil { + p.scheduler = newInprocScheduler(p.coreCtx.Schedules(), p.coreCtx.Tasks()) + } + + return p.scheduler.Start() +} + +// Stop terminates the in-process cron scheduler. +func (p *Plugin) Stop(_ context.Context) error { + p.mu.Lock() + defer p.mu.Unlock() + + if p.scheduler != nil { + p.scheduler.Stop() + p.scheduler = nil + } + return nil +} diff --git a/backend/plugins/drivers/driver_inproc_cron/plugin_test.go b/backend/plugins/drivers/driver_inproc_cron/plugin_test.go new file mode 100644 index 00000000..31e75e62 --- /dev/null +++ b/backend/plugins/drivers/driver_inproc_cron/plugin_test.go @@ -0,0 +1,49 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package driver_inproc_cron_test + +import ( + "context" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "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()) + require.NoError(t, cronPlugin.Apply(ctx)) + + var cronTriggered atomic.Int32 + ctx.Tasks().Register("cron:heartbeat", func(ctx context.Context) error { + cronTriggered.Add(1) + return nil + }) + + // Every second cron spec + ctx.Schedules().RegisterCron("* * * * * *", "cron:heartbeat", nil) + + require.NoError(t, cronPlugin.Start(context.Background())) + + require.Eventually(t, func() bool { + return cronTriggered.Load() >= 1 + }, 3*time.Second, 100*time.Millisecond) + + require.NoError(t, cronPlugin.Stop(context.Background())) +} diff --git a/backend/plugins/drivers/driver_inproc_cron/scheduler.go b/backend/plugins/drivers/driver_inproc_cron/scheduler.go new file mode 100644 index 00000000..6a0fc509 --- /dev/null +++ b/backend/plugins/drivers/driver_inproc_cron/scheduler.go @@ -0,0 +1,116 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package driver_inproc_cron + +import ( + "context" + "encoding/json" + "sync" + + "Wavelet/core/extpoints" + "Wavelet/pkg/logger" + "Wavelet/plugins/drivers/driver_inproc_worker" + "github.com/robfig/cron/v3" +) + +type inprocScheduler struct { + mu sync.RWMutex + cronRunner *cron.Cron + scheduleReg extpoints.ScheduleExtension + taskReg extpoints.TaskExtension + running bool +} + +func newInprocScheduler(scheduleReg extpoints.ScheduleExtension, taskReg extpoints.TaskExtension) *inprocScheduler { + return &inprocScheduler{ + cronRunner: cron.New(cron.WithSeconds()), + scheduleReg: scheduleReg, + taskReg: taskReg, + } +} + +func (s *inprocScheduler) Start() error { + s.mu.Lock() + defer s.mu.Unlock() + + if s.running { + return nil + } + + if s.scheduleReg != nil { + for _, def := range s.scheduleReg.Schedules() { + s.registerJob(def) + } + } + + s.cronRunner.Start() + s.running = true + return nil +} + +func (s *inprocScheduler) Stop() { + s.mu.Lock() + defer s.mu.Unlock() + + if !s.running { + return + } + + ctx := s.cronRunner.Stop() + <-ctx.Done() + s.running = false +} + +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 { + cronSpec = "0 " + spec + } + + var payloadBytes []byte + if def.Payload != nil { + switch p := def.Payload.(type) { + case []byte: + payloadBytes = p + case string: + payloadBytes = []byte(p) + default: + payloadBytes, _ = json.Marshal(p) + } + } + + _, 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 err != nil { + logger.ErrorF(context.Background(), "driver_inproc_cron: invalid cron spec %q for task %q: %v", spec, taskType, err) + } +} + +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 +}