perf(cron): index jobs by id for lower overhead
This commit is contained in:
parent
83e24e8ceb
commit
3c637889d6
2 changed files with 130 additions and 53 deletions
|
|
@ -61,6 +61,7 @@ type JobHandler func(job *CronJob) (string, error)
|
||||||
type CronService struct {
|
type CronService struct {
|
||||||
storePath string
|
storePath string
|
||||||
store *CronStore
|
store *CronStore
|
||||||
|
jobIndex map[string]int
|
||||||
onJob JobHandler
|
onJob JobHandler
|
||||||
mu sync.RWMutex
|
mu sync.RWMutex
|
||||||
running bool
|
running bool
|
||||||
|
|
@ -179,13 +180,9 @@ func (cs *CronService) executeJobByID(jobID string) {
|
||||||
|
|
||||||
cs.mu.RLock()
|
cs.mu.RLock()
|
||||||
var callbackJob *CronJob
|
var callbackJob *CronJob
|
||||||
for i := range cs.store.Jobs {
|
if idx, ok := cs.jobIndex[jobID]; ok && idx >= 0 && idx < len(cs.store.Jobs) && cs.store.Jobs[idx].ID == jobID {
|
||||||
job := &cs.store.Jobs[i]
|
jobCopy := cs.store.Jobs[idx]
|
||||||
if job.ID == jobID {
|
|
||||||
jobCopy := *job
|
|
||||||
callbackJob = &jobCopy
|
callbackJob = &jobCopy
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
cs.mu.RUnlock()
|
cs.mu.RUnlock()
|
||||||
|
|
||||||
|
|
@ -209,14 +206,8 @@ func (cs *CronService) executeJobByID(jobID string) {
|
||||||
cs.mu.Lock()
|
cs.mu.Lock()
|
||||||
defer cs.mu.Unlock()
|
defer cs.mu.Unlock()
|
||||||
|
|
||||||
var job *CronJob
|
job, ok := cs.getJobUnsafe(jobID)
|
||||||
for i := range cs.store.Jobs {
|
if !ok {
|
||||||
if cs.store.Jobs[i].ID == jobID {
|
|
||||||
job = &cs.store.Jobs[i]
|
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if job == nil {
|
|
||||||
log.Printf("[cron] job %s disappeared before state update", jobID)
|
log.Printf("[cron] job %s disappeared before state update", jobID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -338,6 +329,7 @@ func (cs *CronService) loadStore() error {
|
||||||
Version: 1,
|
Version: 1,
|
||||||
Jobs: []CronJob{},
|
Jobs: []CronJob{},
|
||||||
}
|
}
|
||||||
|
cs.jobIndex = make(map[string]int)
|
||||||
|
|
||||||
data, err := os.ReadFile(cs.storePath)
|
data, err := os.ReadFile(cs.storePath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -347,7 +339,34 @@ func (cs *CronService) loadStore() error {
|
||||||
return err
|
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 {
|
func (cs *CronService) saveStoreUnsafe() error {
|
||||||
|
|
@ -396,6 +415,7 @@ func (cs *CronService) AddJob(
|
||||||
}
|
}
|
||||||
|
|
||||||
cs.store.Jobs = append(cs.store.Jobs, job)
|
cs.store.Jobs = append(cs.store.Jobs, job)
|
||||||
|
cs.jobIndex[job.ID] = len(cs.store.Jobs) - 1
|
||||||
if err := cs.saveStoreUnsafe(); err != nil {
|
if err := cs.saveStoreUnsafe(); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
@ -407,16 +427,16 @@ func (cs *CronService) UpdateJob(job *CronJob) error {
|
||||||
cs.mu.Lock()
|
cs.mu.Lock()
|
||||||
defer cs.mu.Unlock()
|
defer cs.mu.Unlock()
|
||||||
|
|
||||||
for i := range cs.store.Jobs {
|
storedJob, ok := cs.getJobUnsafe(job.ID)
|
||||||
if cs.store.Jobs[i].ID == job.ID {
|
if !ok {
|
||||||
cs.store.Jobs[i] = *job
|
|
||||||
cs.store.Jobs[i].UpdatedAtMS = time.Now().UnixMilli()
|
|
||||||
return cs.saveStoreUnsafe()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
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 {
|
func (cs *CronService) RemoveJob(jobID string) bool {
|
||||||
cs.mu.Lock()
|
cs.mu.Lock()
|
||||||
defer cs.mu.Unlock()
|
defer cs.mu.Unlock()
|
||||||
|
|
@ -425,32 +445,39 @@ func (cs *CronService) RemoveJob(jobID string) bool {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (cs *CronService) removeJobUnsafe(jobID string) bool {
|
func (cs *CronService) removeJobUnsafe(jobID string) bool {
|
||||||
before := len(cs.store.Jobs)
|
idx, ok := cs.jobIndex[jobID]
|
||||||
var jobs []CronJob
|
if !ok || idx < 0 || idx >= len(cs.store.Jobs) {
|
||||||
for _, job := range cs.store.Jobs {
|
return false
|
||||||
if job.ID != jobID {
|
|
||||||
jobs = append(jobs, job)
|
|
||||||
}
|
}
|
||||||
}
|
|
||||||
cs.store.Jobs = jobs
|
|
||||||
removed := len(cs.store.Jobs) < before
|
|
||||||
|
|
||||||
if removed {
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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 {
|
if err := cs.saveStoreUnsafe(); err != nil {
|
||||||
log.Printf("[cron] failed to save store after remove: %v", err)
|
log.Printf("[cron] failed to save store after remove: %v", err)
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
return removed
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (cs *CronService) EnableJob(jobID string, enabled bool) *CronJob {
|
func (cs *CronService) EnableJob(jobID string, enabled bool) *CronJob {
|
||||||
cs.mu.Lock()
|
cs.mu.Lock()
|
||||||
defer cs.mu.Unlock()
|
defer cs.mu.Unlock()
|
||||||
|
|
||||||
for i := range cs.store.Jobs {
|
job, ok := cs.getJobUnsafe(jobID)
|
||||||
job := &cs.store.Jobs[i]
|
if !ok {
|
||||||
if job.ID == jobID {
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
job.Enabled = enabled
|
job.Enabled = enabled
|
||||||
job.UpdatedAtMS = time.Now().UnixMilli()
|
job.UpdatedAtMS = time.Now().UnixMilli()
|
||||||
|
|
||||||
|
|
@ -463,11 +490,8 @@ func (cs *CronService) EnableJob(jobID string, enabled bool) *CronJob {
|
||||||
if err := cs.saveStoreUnsafe(); err != nil {
|
if err := cs.saveStoreUnsafe(); err != nil {
|
||||||
log.Printf("[cron] failed to save store after enable: %v", err)
|
log.Printf("[cron] failed to save store after enable: %v", err)
|
||||||
}
|
}
|
||||||
return job
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
return job
|
||||||
}
|
}
|
||||||
|
|
||||||
func (cs *CronService) ListJobs(includeDisabled bool) []CronJob {
|
func (cs *CronService) ListJobs(includeDisabled bool) []CronJob {
|
||||||
|
|
|
||||||
|
|
@ -36,3 +36,56 @@ func TestSaveStore_FilePermissions(t *testing.T) {
|
||||||
func int64Ptr(v int64) *int64 {
|
func int64Ptr(v int64) *int64 {
|
||||||
return &v
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue