diff --git a/.gitignore b/.gitignore index 61fe494ca..000d72dc6 100644 --- a/.gitignore +++ b/.gitignore @@ -25,6 +25,7 @@ build/ # Secrets & Config (keep templates, ignore actual secrets) .env config/config.json +docker/data/ # Test coverage.txt diff --git a/pkg/cron/service.go b/pkg/cron/service.go index 04775ac42..33dfa64a1 100644 --- a/pkg/cron/service.go +++ b/pkg/cron/service.go @@ -61,6 +61,7 @@ type JobHandler func(job *CronJob) (string, error) type CronService struct { storePath string store *CronStore + jobIndex map[string]int onJob JobHandler mu sync.RWMutex running bool @@ -179,13 +180,9 @@ func (cs *CronService) executeJobByID(jobID string) { cs.mu.RLock() var callbackJob *CronJob - for i := range cs.store.Jobs { - job := &cs.store.Jobs[i] - if job.ID == jobID { - jobCopy := *job - callbackJob = &jobCopy - break - } + if idx, ok := cs.jobIndex[jobID]; ok && idx >= 0 && idx < len(cs.store.Jobs) && cs.store.Jobs[idx].ID == jobID { + jobCopy := cs.store.Jobs[idx] + callbackJob = &jobCopy } cs.mu.RUnlock() @@ -209,14 +206,8 @@ func (cs *CronService) executeJobByID(jobID string) { cs.mu.Lock() defer cs.mu.Unlock() - var job *CronJob - for i := range cs.store.Jobs { - if cs.store.Jobs[i].ID == jobID { - job = &cs.store.Jobs[i] - break - } - } - if job == nil { + job, ok := cs.getJobUnsafe(jobID) + if !ok { log.Printf("[cron] job %s disappeared before state update", jobID) return } @@ -338,6 +329,7 @@ func (cs *CronService) loadStore() error { Version: 1, Jobs: []CronJob{}, } + cs.jobIndex = make(map[string]int) data, err := os.ReadFile(cs.storePath) if err != nil { @@ -347,7 +339,34 @@ func (cs *CronService) loadStore() error { return err } - return json.Unmarshal(data, cs.store) + if err := json.Unmarshal(data, cs.store); err != nil { + return err + } + + cs.rebuildJobIndexUnsafe() + return nil +} + +func (cs *CronService) rebuildJobIndexUnsafe() { + cs.jobIndex = make(map[string]int, len(cs.store.Jobs)) + for i := range cs.store.Jobs { + cs.jobIndex[cs.store.Jobs[i].ID] = i + } +} + +func (cs *CronService) getJobUnsafe(jobID string) (*CronJob, bool) { + idx, ok := cs.jobIndex[jobID] + if !ok || idx < 0 || idx >= len(cs.store.Jobs) { + return nil, false + } + if cs.store.Jobs[idx].ID != jobID { + cs.rebuildJobIndexUnsafe() + idx, ok = cs.jobIndex[jobID] + if !ok || idx < 0 || idx >= len(cs.store.Jobs) { + return nil, false + } + } + return &cs.store.Jobs[idx], true } func (cs *CronService) saveStoreUnsafe() error { @@ -396,6 +415,7 @@ func (cs *CronService) AddJob( } cs.store.Jobs = append(cs.store.Jobs, job) + cs.jobIndex[job.ID] = len(cs.store.Jobs) - 1 if err := cs.saveStoreUnsafe(); err != nil { return nil, err } @@ -407,14 +427,14 @@ func (cs *CronService) UpdateJob(job *CronJob) error { cs.mu.Lock() defer cs.mu.Unlock() - for i := range cs.store.Jobs { - if cs.store.Jobs[i].ID == job.ID { - cs.store.Jobs[i] = *job - cs.store.Jobs[i].UpdatedAtMS = time.Now().UnixMilli() - return cs.saveStoreUnsafe() - } + storedJob, ok := cs.getJobUnsafe(job.ID) + if !ok { + return fmt.Errorf("job not found") } - return fmt.Errorf("job not found") + + *storedJob = *job + storedJob.UpdatedAtMS = time.Now().UnixMilli() + return cs.saveStoreUnsafe() } func (cs *CronService) RemoveJob(jobID string) bool { @@ -425,49 +445,53 @@ func (cs *CronService) RemoveJob(jobID string) bool { } func (cs *CronService) removeJobUnsafe(jobID string) bool { - before := len(cs.store.Jobs) - var jobs []CronJob - for _, job := range cs.store.Jobs { - if job.ID != jobID { - jobs = append(jobs, job) - } + idx, ok := cs.jobIndex[jobID] + if !ok || idx < 0 || idx >= len(cs.store.Jobs) { + return false } - cs.store.Jobs = jobs - removed := len(cs.store.Jobs) < before - if removed { - if err := cs.saveStoreUnsafe(); err != nil { - log.Printf("[cron] failed to save store after remove: %v", err) + if cs.store.Jobs[idx].ID != jobID { + cs.rebuildJobIndexUnsafe() + idx, ok = cs.jobIndex[jobID] + if !ok || idx < 0 || idx >= len(cs.store.Jobs) { + return false } } - return removed + copy(cs.store.Jobs[idx:], cs.store.Jobs[idx+1:]) + cs.store.Jobs = cs.store.Jobs[:len(cs.store.Jobs)-1] + cs.rebuildJobIndexUnsafe() + + if err := cs.saveStoreUnsafe(); err != nil { + log.Printf("[cron] failed to save store after remove: %v", err) + } + + return true } func (cs *CronService) EnableJob(jobID string, enabled bool) *CronJob { cs.mu.Lock() defer cs.mu.Unlock() - for i := range cs.store.Jobs { - job := &cs.store.Jobs[i] - if job.ID == jobID { - job.Enabled = enabled - job.UpdatedAtMS = time.Now().UnixMilli() - - if enabled { - job.State.NextRunAtMS = cs.computeNextRun(&job.Schedule, time.Now().UnixMilli()) - } else { - job.State.NextRunAtMS = nil - } - - if err := cs.saveStoreUnsafe(); err != nil { - log.Printf("[cron] failed to save store after enable: %v", err) - } - return job - } + job, ok := cs.getJobUnsafe(jobID) + if !ok { + return nil } - return nil + job.Enabled = enabled + job.UpdatedAtMS = time.Now().UnixMilli() + + if enabled { + job.State.NextRunAtMS = cs.computeNextRun(&job.Schedule, time.Now().UnixMilli()) + } else { + job.State.NextRunAtMS = nil + } + + if err := cs.saveStoreUnsafe(); err != nil { + log.Printf("[cron] failed to save store after enable: %v", err) + } + + return job } func (cs *CronService) ListJobs(includeDisabled bool) []CronJob { diff --git a/pkg/cron/service_test.go b/pkg/cron/service_test.go index 1a0dd1829..82a923cf6 100644 --- a/pkg/cron/service_test.go +++ b/pkg/cron/service_test.go @@ -36,3 +36,56 @@ func TestSaveStore_FilePermissions(t *testing.T) { func int64Ptr(v int64) *int64 { return &v } + +func TestCronServiceIndexedOperationsRemainCompatible(t *testing.T) { + tmpDir := t.TempDir() + storePath := filepath.Join(tmpDir, "cron", "jobs.json") + + cs := NewCronService(storePath, nil) + + job1, err := cs.AddJob("job-1", CronSchedule{Kind: "every", EveryMS: int64Ptr(60000)}, "hello", false, "cli", "direct") + if err != nil { + t.Fatalf("AddJob job1 failed: %v", err) + } + job2, err := cs.AddJob("job-2", CronSchedule{Kind: "every", EveryMS: int64Ptr(120000)}, "world", false, "cli", "direct") + if err != nil { + t.Fatalf("AddJob job2 failed: %v", err) + } + + job2.Name = "job-2-updated" + if err := cs.UpdateJob(job2); err != nil { + t.Fatalf("UpdateJob failed: %v", err) + } + + disabled := cs.EnableJob(job1.ID, false) + if disabled == nil { + t.Fatal("EnableJob returned nil for existing job") + } + if disabled.Enabled { + t.Fatal("expected job to be disabled") + } + + if !cs.RemoveJob(job1.ID) { + t.Fatal("RemoveJob failed for existing job") + } + + reloaded := NewCronService(storePath, nil) + jobs := reloaded.ListJobs(true) + if len(jobs) != 1 { + t.Fatalf("expected 1 job after reload, got %d", len(jobs)) + } + if jobs[0].ID != job2.ID { + t.Fatalf("expected remaining job %s, got %s", job2.ID, jobs[0].ID) + } + if jobs[0].Name != "job-2-updated" { + t.Fatalf("expected updated job name to persist, got %q", jobs[0].Name) + } + + enabled := reloaded.EnableJob(job2.ID, true) + if enabled == nil { + t.Fatal("EnableJob on reloaded service returned nil") + } + if !enabled.Enabled { + t.Fatal("expected reloaded job to be enabled") + } +}