mengstack-api/internal/kernel/jobs/scheduler.go
MengStack Dev df809f1045
Some checks failed
CI / Build & Test (push) Failing after 1m31s
feat: user CRUD API + dashboard stats + kernel infrastructure
- Add user management endpoints (list/create/get/update/delete) with pagination and search
- Add dashboard stats endpoint with tenant/user/online counts and growth metrics
- Add tenant resolver middleware for multi-tenant request scoping
- Add i18n kernel with zh/en message files and AcceptLanguage middleware
- Add WebSocket hub/handler for real-time communication
- Add job scheduler kernel with cron support
- Add plugin sandbox for isolated execution
- Add storage kernel (local filesystem)
- Add event bus kernel for pub/sub
- Add cache kernel abstraction
- Add database migration runner and version upgrade checker
- Add rate limiting middleware with Redis backend
- Add SQL migrations for rbac, audit_logs, settings, notifications, examples
- Extend user repository with list/delete/count operations
- Register all module routes with tenant resolver
2026-10-03 03:42:58 +08:00

101 lines
2.4 KiB
Go

package jobs
import (
"context"
"fmt"
"time"
"github.com/go-co-op/gocron/v2"
"go.uber.org/zap"
)
type Job struct {
Name string
Interval time.Duration
Fn func(ctx context.Context) error
}
type CronJob struct {
Name string
CronExpr string
Fn func(ctx context.Context) error
}
type Scheduler struct {
s gocron.Scheduler
log *zap.Logger
}
func NewScheduler(log *zap.Logger) (*Scheduler, error) {
s, err := gocron.NewScheduler()
if err != nil {
return nil, fmt.Errorf("create scheduler: %w", err)
}
return &Scheduler{s: s, log: log}, nil
}
func (s *Scheduler) AddIntervalJob(job Job) error {
_, err := s.s.NewJob(
gocron.DurationJob(job.Interval),
gocron.NewTask(func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
if err := job.Fn(ctx); err != nil {
s.log.Error("job failed", zap.String("name", job.Name), zap.Error(err))
}
}),
gocron.WithName(job.Name),
)
if err != nil {
return fmt.Errorf("add interval job %s: %w", job.Name, err)
}
s.log.Info("registered interval job", zap.String("name", job.Name), zap.Duration("interval", job.Interval))
return nil
}
func (s *Scheduler) AddCronJob(job CronJob) error {
_, err := s.s.NewJob(
gocron.CronJob(job.CronExpr, false),
gocron.NewTask(func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
if err := job.Fn(ctx); err != nil {
s.log.Error("job failed", zap.String("name", job.Name), zap.Error(err))
}
}),
gocron.WithName(job.Name),
)
if err != nil {
return fmt.Errorf("add cron job %s: %w", job.Name, err)
}
s.log.Info("registered cron job", zap.String("name", job.Name), zap.String("cron", job.CronExpr))
return nil
}
func (s *Scheduler) AddOneShotJob(name string, fn func(ctx context.Context) error, delay time.Duration) error {
_, err := s.s.NewJob(
gocron.OneTimeJob(gocron.OneTimeJobStartDateTime(time.Now().Add(delay))),
gocron.NewTask(func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
if err := fn(ctx); err != nil {
s.log.Error("one-shot job failed", zap.String("name", name), zap.Error(err))
}
}),
gocron.WithName(name),
)
if err != nil {
return fmt.Errorf("add one-shot job %s: %w", name, err)
}
return nil
}
func (s *Scheduler) Start() {
s.s.Start()
s.log.Info("job scheduler started")
}
func (s *Scheduler) Stop() error {
return s.s.Shutdown()
}