This commit is contained in:
OpenClaw-User 2026-03-14 17:01:24 +08:00
commit 563ce3682d
3 changed files with 131 additions and 53 deletions

1
.gitignore vendored
View file

@ -25,6 +25,7 @@ build/
# Secrets & Config (keep templates, ignore actual secrets) # Secrets & Config (keep templates, ignore actual secrets)
.env .env
config/config.json config/config.json
docker/data/
# Test # Test
coverage.txt coverage.txt

View file

@ -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,14 +427,14 @@ 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 {
@ -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 {

View file

@ -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")
}
}