2024-05-03 17:11:56 +00:00
|
|
|
package repo
|
|
|
|
|
|
|
|
import (
|
|
|
|
"fmt"
|
2024-05-13 11:40:01 +00:00
|
|
|
"log/slog"
|
2024-05-03 17:11:56 +00:00
|
|
|
"math"
|
2024-05-08 10:24:34 +00:00
|
|
|
"slices"
|
|
|
|
"strings"
|
2024-05-03 17:11:56 +00:00
|
|
|
"time"
|
2024-05-08 14:35:04 +00:00
|
|
|
"xmr-remote-nodes/internal/database"
|
2024-05-03 17:11:56 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
type CronRepository interface {
|
|
|
|
RunCronProcess()
|
2024-05-08 10:24:34 +00:00
|
|
|
Crons(q CronQueryParams) (CronTasks, error)
|
2024-05-03 17:11:56 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
type CronRepo struct {
|
|
|
|
db *database.DB
|
|
|
|
}
|
|
|
|
|
2024-05-08 10:24:34 +00:00
|
|
|
type Cron struct {
|
2024-05-03 17:11:56 +00:00
|
|
|
Id int `json:"id" db:"id"`
|
|
|
|
Title string `json:"title" db:"title"`
|
|
|
|
Slug string `json:"slug" db:"slug"`
|
|
|
|
Description string `json:"description" db:"description"`
|
|
|
|
RunEvery int `json:"run_every" db:"run_every"`
|
|
|
|
LastRun int64 `json:"last_run" db:"last_run"`
|
|
|
|
NextRun int64 `json:"next_run" db:"next_run"`
|
|
|
|
RunTime float64 `json:"run_time" db:"run_time"`
|
|
|
|
CronState int `json:"cron_state" db:"cron_state"`
|
|
|
|
IsEnabled int `json:"is_enabled" db:"is_enabled"`
|
|
|
|
}
|
|
|
|
|
|
|
|
var rerunTimeout = 300
|
|
|
|
|
|
|
|
func NewCron(db *database.DB) CronRepository {
|
|
|
|
return &CronRepo{db}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (repo *CronRepo) RunCronProcess() {
|
|
|
|
for {
|
|
|
|
time.Sleep(60 * time.Second)
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Info("[CRON] Running cron cycle...")
|
2024-05-03 17:11:56 +00:00
|
|
|
list, err := repo.queueList()
|
|
|
|
if err != nil {
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Warn(fmt.Sprintf("[CRON] Error parsing queue list to struct: %s", err))
|
2024-05-03 17:11:56 +00:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
for _, task := range list {
|
|
|
|
startTime := time.Now()
|
|
|
|
currentTs := startTime.Unix()
|
|
|
|
delayedTask := currentTs - task.NextRun
|
|
|
|
if task.CronState == 1 && delayedTask <= int64(rerunTimeout) {
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Debug(fmt.Sprintf("[CRON] Skipping task %s because it is already running", task.Slug))
|
2024-05-03 17:11:56 +00:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
|
|
|
|
repo.preRunTask(task.Id, currentTs)
|
|
|
|
|
|
|
|
repo.execCron(task.Slug)
|
|
|
|
|
|
|
|
runTime := math.Ceil(time.Since(startTime).Seconds()*1000) / 1000
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Info(fmt.Sprintf("[CRON] Task %s done in %f seconds", task.Slug, runTime))
|
2024-05-03 17:11:56 +00:00
|
|
|
nextRun := currentTs + int64(task.RunEvery)
|
|
|
|
|
|
|
|
repo.postRunTask(task.Id, nextRun, runTime)
|
|
|
|
}
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Info("[CRON] Cron cycle done!")
|
2024-05-03 17:11:56 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-05-08 10:24:34 +00:00
|
|
|
type CronQueryParams struct {
|
|
|
|
Title string
|
|
|
|
Description string
|
|
|
|
IsEnabled int
|
|
|
|
CronState int
|
|
|
|
RowsPerPage int
|
|
|
|
Page int
|
|
|
|
SortBy string
|
|
|
|
SortDirection string
|
|
|
|
}
|
|
|
|
|
|
|
|
type CronTasks struct {
|
|
|
|
TotalRows int `json:"total_rows"`
|
|
|
|
RowsPerPage int `json:"rows_per_page"`
|
|
|
|
Items []*Cron `json:"items"`
|
|
|
|
}
|
|
|
|
|
|
|
|
func (repo *CronRepo) Crons(q CronQueryParams) (CronTasks, error) {
|
|
|
|
queryParams := []interface{}{}
|
|
|
|
whereQueries := []string{}
|
|
|
|
where := ""
|
|
|
|
|
|
|
|
if q.Title != "" {
|
|
|
|
whereQueries = append(whereQueries, "title LIKE ?")
|
|
|
|
queryParams = append(queryParams, "%"+q.Title+"%")
|
|
|
|
}
|
|
|
|
if q.Description != "" {
|
|
|
|
whereQueries = append(whereQueries, "description LIKE ?")
|
|
|
|
queryParams = append(queryParams, "%"+q.Description+"%")
|
|
|
|
}
|
|
|
|
if q.IsEnabled != -1 {
|
|
|
|
whereQueries = append(whereQueries, "is_enabled = ?")
|
|
|
|
queryParams = append(queryParams, q.IsEnabled)
|
|
|
|
}
|
|
|
|
if q.CronState != -1 {
|
|
|
|
whereQueries = append(whereQueries, "cron_state = ?")
|
|
|
|
queryParams = append(queryParams, q.CronState)
|
|
|
|
}
|
|
|
|
if len(whereQueries) > 0 {
|
|
|
|
where = "WHERE " + strings.Join(whereQueries, " AND ")
|
|
|
|
}
|
|
|
|
tasks := CronTasks{}
|
|
|
|
|
|
|
|
queryTotalRows := fmt.Sprintf("SELECT COUNT(id) FROM tbl_cron %s", where)
|
|
|
|
err := repo.db.QueryRow(queryTotalRows, queryParams...).Scan(&tasks.TotalRows)
|
|
|
|
if err != nil {
|
|
|
|
return tasks, err
|
|
|
|
}
|
|
|
|
queryParams = append(queryParams, q.RowsPerPage, (q.Page-1)*q.RowsPerPage)
|
|
|
|
allowedSort := []string{"id", "run_every", "last_run", "next_run", "run_time"}
|
|
|
|
sortBy := "id"
|
|
|
|
if slices.Contains(allowedSort, q.SortBy) {
|
|
|
|
sortBy = q.SortBy
|
|
|
|
}
|
|
|
|
sortDirection := "DESC"
|
|
|
|
if q.SortDirection == "asc" {
|
|
|
|
sortDirection = "ASC"
|
|
|
|
}
|
|
|
|
|
|
|
|
query := fmt.Sprintf("SELECT id, title, slug, description, run_every, last_run, next_run, run_time, cron_state, is_enabled FROM tbl_cron %s ORDER BY %s %s LIMIT ? OFFSET ?", where, sortBy, sortDirection)
|
|
|
|
err = repo.db.Select(&tasks.Items, query, queryParams...)
|
|
|
|
if err != nil {
|
|
|
|
return tasks, err
|
|
|
|
}
|
|
|
|
tasks.RowsPerPage = q.RowsPerPage
|
|
|
|
|
|
|
|
return tasks, nil
|
2024-05-03 17:11:56 +00:00
|
|
|
}
|
|
|
|
|
2024-05-08 10:24:34 +00:00
|
|
|
func (repo *CronRepo) queueList() ([]Cron, error) {
|
|
|
|
tasks := []Cron{}
|
2024-05-03 17:11:56 +00:00
|
|
|
query := `SELECT id, run_every, last_run, slug, next_run, cron_state FROM tbl_cron
|
|
|
|
WHERE is_enabled = ? AND next_run <= ?`
|
|
|
|
err := repo.db.Select(&tasks, query, 1, time.Now().Unix())
|
|
|
|
|
|
|
|
return tasks, err
|
|
|
|
}
|
|
|
|
|
|
|
|
func (repo *CronRepo) preRunTask(id int, lastRunTs int64) {
|
|
|
|
query := `UPDATE tbl_cron SET cron_state = ?, last_run = ? WHERE id = ?`
|
|
|
|
row, err := repo.db.Query(query, 1, lastRunTs, id)
|
|
|
|
if err != nil {
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Error(fmt.Sprintf("[CRON] Failed to update pre cron state: %s", err))
|
2024-05-03 17:11:56 +00:00
|
|
|
}
|
|
|
|
defer row.Close()
|
|
|
|
}
|
|
|
|
|
|
|
|
func (repo *CronRepo) postRunTask(id int, nextRun int64, runtime float64) {
|
|
|
|
query := `UPDATE tbl_cron SET cron_state = ?, next_run = ?, run_time = ? WHERE id = ?`
|
|
|
|
row, err := repo.db.Query(query, 0, nextRun, runtime, id)
|
|
|
|
if err != nil {
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Error(fmt.Sprintf("[CRON] Failed to update post cron state: %s", err))
|
2024-05-03 17:11:56 +00:00
|
|
|
}
|
|
|
|
defer row.Close()
|
|
|
|
}
|
|
|
|
|
|
|
|
func (repo *CronRepo) execCron(slug string) {
|
|
|
|
switch slug {
|
2024-05-06 11:40:09 +00:00
|
|
|
case "delete_old_probe_logs":
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Info(fmt.Sprintf("[CRON] Start running task: %s", slug))
|
2024-05-06 11:40:09 +00:00
|
|
|
repo.deleteOldProbeLogs()
|
2024-05-03 17:11:56 +00:00
|
|
|
break
|
|
|
|
}
|
|
|
|
}
|
2024-05-06 11:40:09 +00:00
|
|
|
|
|
|
|
func (repo *CronRepo) deleteOldProbeLogs() {
|
2024-05-08 12:28:42 +00:00
|
|
|
// for now, we only delete stats older than 1 month +2 days
|
|
|
|
startTs := time.Now().AddDate(0, -1, -2).Unix()
|
2024-05-06 11:40:09 +00:00
|
|
|
query := `DELETE FROM tbl_probe_log WHERE date_checked < ?`
|
|
|
|
_, err := repo.db.Exec(query, startTs)
|
|
|
|
if err != nil {
|
2024-05-13 11:40:01 +00:00
|
|
|
slog.Error(fmt.Sprintf("[CRON] Failed to delete old probe logs: %s", err))
|
2024-05-06 11:40:09 +00:00
|
|
|
}
|
|
|
|
}
|