mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-01 06:36:38 +08:00
feat(drivers): implement in-process cron scheduler driver plugin
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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()))
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user