58e7f51ba7
- eventbus.go: Fix send-on-closed-channel panic in SubscribeGlobal by using a done channel; add recover() in fan-out goroutine; track active global subs for proper cleanup on unsubscribe - config.go: Persist JWT secret to $DATA_DIR/.jwt_secret instead of regenerating a random one on every restart (which invalidated all sessions) - handlers_ws.go: Replace time.After with time.Ticker to fix timer leak in SSE keepalive loop - handlers_jobs.go: Add recover() in fire-and-forget job goroutine; fix nil pointer deref when GetByID fails after job creation - handlers_machines.go: Add recover() in ProbeAllMachines goroutine - scheduler.go: Add recover() in scheduled job run goroutine - engine.go: Add recover() in per-machine probe goroutines
146 lines
3.1 KiB
Go
146 lines
3.1 KiB
Go
package scheduler
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"log/slog"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/syncserver/internal/config"
|
|
"github.com/syncserver/internal/models"
|
|
"github.com/syncserver/internal/syncengine"
|
|
)
|
|
|
|
type Scheduler struct {
|
|
db *sql.DB
|
|
engine *syncengine.Engine
|
|
cfg *config.Config
|
|
stopCh chan struct{}
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
func New(database interface{ SQLDB() *sql.DB }, engine *syncengine.Engine, cfg *config.Config) *Scheduler {
|
|
return &Scheduler{
|
|
db: database.SQLDB(),
|
|
engine: engine,
|
|
cfg: cfg,
|
|
stopCh: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
func (s *Scheduler) Start() {
|
|
s.wg.Add(1)
|
|
go s.run()
|
|
s.wg.Add(1)
|
|
go s.cleanupRun()
|
|
slog.Info("scheduler started")
|
|
}
|
|
|
|
func (s *Scheduler) Stop() {
|
|
close(s.stopCh)
|
|
s.wg.Wait()
|
|
slog.Info("scheduler stopped")
|
|
}
|
|
|
|
func (s *Scheduler) run() {
|
|
defer s.wg.Done()
|
|
ticker := time.NewTicker(1 * time.Minute)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-s.stopCh:
|
|
return
|
|
case <-ticker.C:
|
|
s.tick()
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Scheduler) tick() {
|
|
scheduleRepo := models.NewScheduleRepository(s.db)
|
|
now := time.Now().UTC()
|
|
|
|
schedules, err := scheduleRepo.GetEnabledDue(now)
|
|
if err != nil {
|
|
slog.Error("scheduler: failed to get due schedules", "error", err)
|
|
return
|
|
}
|
|
|
|
for _, sch := range schedules {
|
|
pairRepo := models.NewSyncPairRepository(s.db)
|
|
pair, err := pairRepo.GetByID(sch.SyncPairID)
|
|
if err != nil || !pair.Enabled {
|
|
continue
|
|
}
|
|
|
|
jobID, err := s.engine.CreateJob(sch.SyncPairID, "scheduled")
|
|
if err != nil {
|
|
slog.Error("scheduler: failed to create job", "schedule_id", sch.ID, "error", err)
|
|
continue
|
|
}
|
|
|
|
ctx := context.Background()
|
|
go func(jobID int64, pairID int64, schID int64) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
slog.Error("scheduler: job run panicked", "job_id", jobID, "panic", r)
|
|
}
|
|
}()
|
|
if err := s.engine.Run(ctx, jobID, pairID); err != nil {
|
|
slog.Warn("scheduler: job failed", "job_id", jobID, "error", err)
|
|
}
|
|
|
|
expr, _ := ParseCron(sch.CronExpr)
|
|
if expr != nil {
|
|
next := NextRun(expr, time.Now().UTC())
|
|
scheduleRepo.UpdateNextRun(schID, next)
|
|
}
|
|
}(jobID, sch.SyncPairID, sch.ID)
|
|
}
|
|
}
|
|
|
|
func (s *Scheduler) cleanupRun() {
|
|
defer s.wg.Done()
|
|
ticker := time.NewTicker(24 * time.Hour)
|
|
defer ticker.Stop()
|
|
|
|
s.cleanup()
|
|
for {
|
|
select {
|
|
case <-s.stopCh:
|
|
return
|
|
case <-ticker.C:
|
|
s.cleanup()
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Scheduler) cleanup() {
|
|
retentionDays := s.cfg.Scheduler.RetentionDays
|
|
if retentionDays <= 0 {
|
|
return
|
|
}
|
|
before := time.Now().AddDate(0, 0, -retentionDays)
|
|
|
|
logRepo := models.NewJobLogRepository(s.db)
|
|
jobRepo := models.NewJobRepository(s.db)
|
|
|
|
deletedLogs, err := logRepo.DeleteBefore(before)
|
|
if err != nil {
|
|
slog.Error("cleanup: failed to purge old job logs", "error", err)
|
|
return
|
|
}
|
|
|
|
deletedJobs, err := jobRepo.DeleteFinishedBefore(before)
|
|
if err != nil {
|
|
slog.Error("cleanup: failed to purge old jobs", "error", err)
|
|
return
|
|
}
|
|
|
|
if deletedLogs > 0 || deletedJobs > 0 {
|
|
slog.Info("cleanup: purged old records", "logs_deleted", deletedLogs, "jobs_deleted", deletedJobs, "before", before.Format("2006-01-02"))
|
|
}
|
|
}
|