Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion main.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ var viewsFS embed.FS
var staticFS embed.FS

var (
Version = "0.14.0"
Version = "0.15.0"
Package = "community"
)

Expand Down
7 changes: 7 additions & 0 deletions models/schedule.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
140 changes: 140 additions & 0 deletions models/schedule_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
11 changes: 11 additions & 0 deletions services/periodic_task_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading