From a2fd4b0381e693f7d828ce0e6cfde13bd7b6a9bd Mon Sep 17 00:00:00 2001 From: Ettore Di Giacinto Date: Mon, 24 Aug 2026 19:34:30 +0000 Subject: [PATCH] fix(scheduler): skip the update on a one-time task it just deleted executeTask deleted a one-time task and then fell through to store.Update on the same row. The update could never find it, so every one-time task logged "Failed to update task: task not found" as it finished. The one-time branch now returns after the delete. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01KZaoEXfGsjmkhXtVxtvPp4 --- core/scheduler/onetime_test.go | 98 ++++++++++++++++++++++++++++++++++ core/scheduler/scheduler.go | 16 +++--- 2 files changed, 107 insertions(+), 7 deletions(-) create mode 100644 core/scheduler/onetime_test.go diff --git a/core/scheduler/onetime_test.go b/core/scheduler/onetime_test.go new file mode 100644 index 0000000..90d481e --- /dev/null +++ b/core/scheduler/onetime_test.go @@ -0,0 +1,98 @@ +package scheduler_test + +import ( + "os" + "sync" + "time" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "github.com/mudler/LocalAGI/core/scheduler" +) + +// recordingStore notes which mutating calls the scheduler makes, so a spec can +// assert on the sequence rather than on a log line. +type recordingStore struct { + scheduler.TaskStore + + mu sync.Mutex + calls []string +} + +func (r *recordingStore) record(call string) { + r.mu.Lock() + defer r.mu.Unlock() + r.calls = append(r.calls, call) +} + +func (r *recordingStore) Calls() []string { + r.mu.Lock() + defer r.mu.Unlock() + return append([]string(nil), r.calls...) +} + +func (r *recordingStore) Delete(id string) error { + r.record("delete:" + id) + return r.TaskStore.Delete(id) +} + +func (r *recordingStore) Update(task *scheduler.Task) error { + r.record("update:" + task.ID) + return r.TaskStore.Update(task) +} + +var _ = Describe("One-time task execution", func() { + var store *recordingStore + var sched *scheduler.Scheduler + + BeforeEach(func() { + f, err := os.CreateTemp("", "onetime_test_*.json") + Expect(err).NotTo(HaveOccurred()) + name := f.Name() + f.Close() + DeferCleanup(func() { os.Remove(name) }) + + base, err := scheduler.NewJSONStore(name) + Expect(err).NotTo(HaveOccurred()) + + store = &recordingStore{TaskStore: base} + sched = scheduler.NewScheduler(store, &MockExecutor{}, 50*time.Millisecond) + sched.Start() + DeferCleanup(sched.Stop) + }) + + // The update used to run unconditionally after the delete, so every + // one-time task logged "task not found" on the way out. + It("does not update a one-time task after deleting it", func() { + task, err := scheduler.NewTask("agent", "ping", scheduler.ScheduleTypeOnce, "0s") + Expect(err).NotTo(HaveOccurred()) + task.NextRun = time.Now().Add(-time.Second) + + _, err = sched.CreateTask(task) + Expect(err).NotTo(HaveOccurred()) + + Eventually(func() []string { + return store.Calls() + }, "3s", "50ms").Should(ContainElement("delete:" + task.ID)) + + Consistently(func() []string { + return store.Calls() + }, "300ms", "50ms").ShouldNot(ContainElement("update:" + task.ID)) + }) + + It("still reschedules a recurring task after it runs", func() { + task, err := scheduler.NewTask("agent", "ping", scheduler.ScheduleTypeCron, "* * * * *") + Expect(err).NotTo(HaveOccurred()) + task.NextRun = time.Now().Add(-time.Second) + + _, err = sched.CreateTask(task) + Expect(err).NotTo(HaveOccurred()) + + Eventually(func() []string { + return store.Calls() + }, "3s", "50ms").Should(ContainElement("update:" + task.ID)) + + Expect(store.Calls()).ToNot(ContainElement("delete:" + task.ID)) + }) +}) diff --git a/core/scheduler/scheduler.go b/core/scheduler/scheduler.go index b273fe5..5ef438a 100644 --- a/core/scheduler/scheduler.go +++ b/core/scheduler/scheduler.go @@ -160,17 +160,19 @@ func (s *Scheduler) executeTask(task *Task) { now := time.Now() task.LastRun = &now - // For one-time tasks, mark as deleted + // A one-time task is done once it has run. Returning here matters: the + // update below would otherwise run against the row just deleted and log a + // "task not found" error on every one-time task. if task.ScheduleType == ScheduleTypeOnce { if err := s.store.Delete(task.ID); err != nil { xlog.Error("Failed to delete task", "task_id", task.ID, "error", err) } - } else { - // Calculate next run - if err := task.CalculateNextRun(); err != nil { - xlog.Error("Failed to calculate next run", "task_id", task.ID, "error", err) - task.Status = TaskStatusPaused - } + return + } + + if err := task.CalculateNextRun(); err != nil { + xlog.Error("Failed to calculate next run", "task_id", task.ID, "error", err) + task.Status = TaskStatusPaused } if err := s.store.Update(task); err != nil {