diff --git a/main.go b/main.go index 73b5d18..74a682d 100644 --- a/main.go +++ b/main.go @@ -33,7 +33,7 @@ var viewsFS embed.FS var staticFS embed.FS var ( - Version = "0.14.0" + Version = "0.15.0" Package = "community" ) diff --git a/models/schedule.go b/models/schedule.go index 484a659..1ee1f4f 100644 --- a/models/schedule.go +++ b/models/schedule.go @@ -51,6 +51,13 @@ func (s *Store) RecordTaskFailure(id int64) error { return err } +func (s *Store) IsPeriodicTaskActive(id int64) (bool, error) { + var active bool + query := `SELECT active FROM periodic_tasks WHERE id = ?` + err := s.db.Get(&active, query, id) + return active, err +} + func (s *Store) GetAndLockDuePeriodicTasks() ([]DueTasksTaskDetail, error) { // TODO: think about this one // err := s.UnlockStalePeriodicTasks() diff --git a/models/schedule_test.go b/models/schedule_test.go new file mode 100644 index 0000000..bcb1428 --- /dev/null +++ b/models/schedule_test.go @@ -0,0 +1,140 @@ +package models + +import ( + "testing" + "time" + + "github.com/jmoiron/sqlx" + "github.com/rapidforge-io/rapidforge/database" + _ "modernc.org/sqlite" +) + +func newScheduleTestStore(t *testing.T) *Store { + t.Helper() + + db := sqlx.MustConnect("sqlite", ":memory:") + t.Cleanup(func() { + db.Close() + }) + + schema := []string{ + `CREATE TABLE blocks ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + env_variables TEXT + )`, + `CREATE TABLE programs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + type TEXT NOT NULL DEFAULT 'bash', + created_at TIMESTAMP + )`, + `CREATE TABLE files ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + program_id INTEGER, + filename TEXT, + content TEXT, + created_at TIMESTAMP + )`, + `CREATE TABLE periodic_tasks ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + description TEXT, + active BOOLEAN DEFAULT 1, + env_variables TEXT, + block_id INTEGER, + program_id INTEGER, + timezone TEXT DEFAULT 'UTC', + cron TEXT NOT NULL, + next_run_at DATETIME, + on_fail_script TEXT, + on_fail_script_type TEXT DEFAULT 'bash', + on_fail_enabled BOOLEAN DEFAULT 0, + locked BOOLEAN DEFAULT 0, + locked_at DATETIME, + created_at TIMESTAMP, + updated_at TIMESTAMP + )`, + } + for _, stmt := range schema { + db.MustExec(stmt) + } + + return &Store{db: &database.DbCon{DB: db}} +} + +func insertScheduleTestTask(t *testing.T, store *Store, active bool, nextRunAt time.Time) int64 { + t.Helper() + + res, err := store.db.Exec(`INSERT INTO blocks (name) VALUES (?)`, "block") + if err != nil { + t.Fatalf("insert block: %v", err) + } + blockID, _ := res.LastInsertId() + + res, err = store.db.Exec(`INSERT INTO programs (name, type, created_at) VALUES (?, ?, ?)`, "program", "bash", time.Now().UTC()) + if err != nil { + t.Fatalf("insert program: %v", err) + } + programID, _ := res.LastInsertId() + + if _, err := store.db.Exec(`INSERT INTO files (program_id, filename, content, created_at) VALUES (?, ?, ?, ?)`, programID, "main", "echo ok", time.Now().UTC()); err != nil { + t.Fatalf("insert file: %v", err) + } + + res, err = store.db.Exec( + `INSERT INTO periodic_tasks (name, active, block_id, program_id, cron, next_run_at) VALUES (?, ?, ?, ?, ?, ?)`, + "task", + active, + blockID, + programID, + "* * * * *", + nextRunAt.UTC(), + ) + if err != nil { + t.Fatalf("insert periodic task: %v", err) + } + + taskID, _ := res.LastInsertId() + return taskID +} + +func TestGetAndLockDuePeriodicTasksSkipsDisabledTasks(t *testing.T) { + store := newScheduleTestStore(t) + dueAt := time.Now().UTC().Add(-time.Minute) + + disabledID := insertScheduleTestTask(t, store, false, dueAt) + activeID := insertScheduleTestTask(t, store, true, dueAt) + + tasks, err := store.GetAndLockDuePeriodicTasks() + if err != nil { + t.Fatalf("GetAndLockDuePeriodicTasks() error = %v", err) + } + if len(tasks) != 1 { + t.Fatalf("GetAndLockDuePeriodicTasks() returned %d tasks, want 1", len(tasks)) + } + if tasks[0].PeriodicTask.ID != activeID { + t.Fatalf("GetAndLockDuePeriodicTasks() returned task %d, want active task %d", tasks[0].PeriodicTask.ID, activeID) + } + + var disabledLocked bool + if err := store.db.Get(&disabledLocked, `SELECT locked FROM periodic_tasks WHERE id = ?`, disabledID); err != nil { + t.Fatalf("select disabled locked state: %v", err) + } + if disabledLocked { + t.Fatal("disabled due task was locked for execution") + } +} + +func TestIsPeriodicTaskActive(t *testing.T) { + store := newScheduleTestStore(t) + taskID := insertScheduleTestTask(t, store, false, time.Now().UTC()) + + active, err := store.IsPeriodicTaskActive(taskID) + if err != nil { + t.Fatalf("IsPeriodicTaskActive() error = %v", err) + } + if active { + t.Fatal("IsPeriodicTaskActive() = true, want false") + } +} diff --git a/services/periodic_task_service.go b/services/periodic_task_service.go index d8788c4..4ead634 100644 --- a/services/periodic_task_service.go +++ b/services/periodic_task_service.go @@ -45,6 +45,17 @@ func (s *Service) RunPeriodicPrograms() { } }() + active, err := s.store.IsPeriodicTaskActive(task.PeriodicTask.ID) + if err != nil { + rflog.Error("Failed to check periodic task active state", "task_id=", task.PeriodicTask.ID, "err=", err) + s.store.UnlockPeriodicTask(task.PeriodicTask.ID) + return + } + if !active { + s.store.UnlockPeriodicTask(task.PeriodicTask.ID) + return + } + blockEnv := task.Block.GetEnvVars() taskEnv := task.PeriodicTask.GetEnvVars() env := utils.MergeMaps(blockEnv, taskEnv)