mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-11 03:13:24 +00:00
fix(cron): wait for in-flight jobs on shutdown
This commit is contained in:
1 parent
8761ce5401
commit
3ec20dfb55
4 files changed
+65
No files matched your search
@@ -6,6 +6,18 @@ Significant changes, features, and fixes in reverse chronological order.
|
||||
|
||||
## 2026-06-12
|
||||
|
||||
### Cron scheduler shutdown drain (post-merge CI follow-up)
|
||||
|
||||
**Fixes**
|
||||
|
||||
- Made cron service shutdown wait for scheduler-launched job executions to
|
||||
finish so in-flight status persistence cannot outlive `Stop()`.
|
||||
|
||||
**Tests**
|
||||
|
||||
- Added regression coverage proving `Stop()` waits for a blocked in-flight job
|
||||
before returning.
|
||||
|
||||
### Bundled GoClaw gateway administration skill (issue #175)
|
||||
|
||||
**Changes**
|
||||
|
||||
@@ -14,6 +14,7 @@ type Service struct {
|
||||
running bool
|
||||
stopChan chan struct{}
|
||||
loopWG sync.WaitGroup
|
||||
jobWG sync.WaitGroup
|
||||
mu sync.Mutex
|
||||
runLog []RunLogEntry // in-memory run history (last 200 entries)
|
||||
retryCfg RetryConfig // retry config for failed jobs
|
||||
@@ -114,6 +115,7 @@ func (cs *Service) Stop() {
|
||||
cs.mu.Unlock()
|
||||
|
||||
cs.loopWG.Wait()
|
||||
cs.jobWG.Wait()
|
||||
slog.Info("cron service stopped")
|
||||
}
|
||||
|
||||
|
||||
@@ -209,6 +209,7 @@ func (cs *Service) checkJobs() {
|
||||
}
|
||||
}
|
||||
cs.saveUnsafe()
|
||||
cs.jobWG.Add(len(dueJobs))
|
||||
cs.mu.Unlock()
|
||||
|
||||
// Execute jobs in parallel without blocking the runLoop.
|
||||
@@ -217,6 +218,7 @@ func (cs *Service) checkJobs() {
|
||||
// due jobs. Now each job runs independently with panic recovery.
|
||||
for _, dj := range dueJobs {
|
||||
go func(id string, scheduledAtMS int64) {
|
||||
defer cs.jobWG.Done()
|
||||
defer safego.Recover(nil, "job_id", id)
|
||||
cs.executeJobByID(id, scheduledAtMS)
|
||||
}(dj.id, dj.scheduledAtMS)
|
||||
|
||||
@@ -306,6 +306,55 @@ func TestService_JobFailure_Updates_LastError(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestService_Stop_WaitsForInFlightJob(t *testing.T) {
|
||||
setFastTick(t)
|
||||
dir := t.TempDir()
|
||||
storePath := filepath.Join(dir, "cron.json")
|
||||
|
||||
entered := make(chan struct{})
|
||||
release := make(chan struct{})
|
||||
handler := func(job *Job) (string, error) {
|
||||
close(entered)
|
||||
<-release
|
||||
return "done", nil
|
||||
}
|
||||
|
||||
cs := NewService(storePath, handler)
|
||||
interval := int64(10)
|
||||
if _, err := cs.AddJob("blocking", Schedule{Kind: "every", EveryMS: &interval}, "block", false, "", "", ""); err != nil {
|
||||
t.Fatalf("AddJob error: %v", err)
|
||||
}
|
||||
if err := cs.Start(); err != nil {
|
||||
t.Fatalf("Start error: %v", err)
|
||||
}
|
||||
|
||||
select {
|
||||
case <-entered:
|
||||
case <-time.After(time.Second):
|
||||
cs.Stop()
|
||||
t.Fatal("job handler did not start")
|
||||
}
|
||||
|
||||
stopped := make(chan struct{})
|
||||
go func() {
|
||||
cs.Stop()
|
||||
close(stopped)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-stopped:
|
||||
t.Fatal("Stop returned while job handler was still running")
|
||||
case <-time.After(50 * time.Millisecond):
|
||||
}
|
||||
|
||||
close(release)
|
||||
select {
|
||||
case <-stopped:
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("Stop did not return after in-flight job completed")
|
||||
}
|
||||
}
|
||||
|
||||
// --- Persistence: save and reload ---
|
||||
|
||||
func TestService_Persistence_Roundtrip(t *testing.T) {
|
||||
|
||||
Reference in new issue
Block a user