Some checks failed
CI Docker Mining Proof / Linux agent hashrate proof (push) Has been cancelled
Fix macOS agent cross-compile (SilentAVExclusion) and Calibrate E2E nav selector; expand tests and docs; refresh portable usb binary and spread/wiki assets.
136 lines
4.2 KiB
Go
136 lines
4.2 KiB
Go
package db
|
||
|
||
import (
|
||
"database/sql"
|
||
"fmt"
|
||
"strings"
|
||
"time"
|
||
|
||
"crypto-miner-server/internal/models"
|
||
|
||
"github.com/google/uuid"
|
||
)
|
||
|
||
func (d *Database) ListFleetTasks() ([]*models.FleetTask, error) {
|
||
rows, err := d.Query(`SELECT id, name, enabled, trigger, interval_hours, cron_time, action, command, target, created_at, updated_at FROM fleet_tasks ORDER BY created_at`)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
return scanFleetTasks(rows)
|
||
}
|
||
|
||
func (d *Database) GetFleetTask(id string) (*models.FleetTask, error) {
|
||
row := d.QueryRow(`SELECT id, name, enabled, trigger, interval_hours, cron_time, action, command, target, created_at, updated_at FROM fleet_tasks WHERE id = ?`, id)
|
||
return scanFleetTaskRow(row)
|
||
}
|
||
|
||
func (d *Database) UpsertFleetTask(t *models.FleetTask) error {
|
||
if t.ID == "" {
|
||
t.ID = uuid.New().String()
|
||
}
|
||
now := time.Now()
|
||
if t.CreatedAt.IsZero() {
|
||
t.CreatedAt = now
|
||
}
|
||
t.UpdatedAt = now
|
||
_, err := d.Exec(`INSERT INTO fleet_tasks (id, name, enabled, trigger, interval_hours, cron_time, action, command, target, created_at, updated_at)
|
||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||
ON CONFLICT(id) DO UPDATE SET
|
||
name = excluded.name,
|
||
enabled = excluded.enabled,
|
||
trigger = excluded.trigger,
|
||
interval_hours = excluded.interval_hours,
|
||
cron_time = excluded.cron_time,
|
||
action = excluded.action,
|
||
command = excluded.command,
|
||
target = excluded.target,
|
||
updated_at = excluded.updated_at`,
|
||
t.ID, t.Name, boolToInt(t.Enabled), t.Trigger, t.IntervalHours, t.CronTime, t.Action, t.Command, t.Target, t.CreatedAt, t.UpdatedAt,
|
||
)
|
||
return err
|
||
}
|
||
|
||
func (d *Database) DeleteFleetTask(id string) error {
|
||
_, err := d.Exec(`DELETE FROM fleet_tasks WHERE id = ?`, id)
|
||
return err
|
||
}
|
||
|
||
func (d *Database) RecordFleetTaskRun(agentID, taskID string) error {
|
||
_, err := d.Exec(`INSERT INTO fleet_task_runs (agent_id, task_id, last_run_at) VALUES (?, ?, ?)
|
||
ON CONFLICT(agent_id, task_id) DO UPDATE SET last_run_at = excluded.last_run_at`,
|
||
agentID, taskID, time.Now(),
|
||
)
|
||
return err
|
||
}
|
||
|
||
func (d *Database) LastFleetTaskRun(agentID, taskID string) (time.Time, bool) {
|
||
var ts time.Time
|
||
err := d.QueryRow(`SELECT last_run_at FROM fleet_task_runs WHERE agent_id = ? AND task_id = ?`, agentID, taskID).Scan(&ts)
|
||
if err != nil {
|
||
return time.Time{}, false
|
||
}
|
||
return ts, true
|
||
}
|
||
|
||
// BulkLastFleetTaskRuns returns a map keyed by "agentID:taskID" with the
|
||
// last_run_at time for every matching row. Missing pairs were never run.
|
||
// A single query replaces O(tasks × agents) individual lookups.
|
||
func (d *Database) BulkLastFleetTaskRuns(agentIDs, taskIDs []string) (map[string]time.Time, error) {
|
||
if len(agentIDs) == 0 || len(taskIDs) == 0 {
|
||
return map[string]time.Time{}, nil
|
||
}
|
||
args := make([]interface{}, 0, len(agentIDs)+len(taskIDs))
|
||
for _, id := range agentIDs {
|
||
args = append(args, id)
|
||
}
|
||
for _, id := range taskIDs {
|
||
args = append(args, id)
|
||
}
|
||
query := fmt.Sprintf(
|
||
`SELECT agent_id, task_id, last_run_at FROM fleet_task_runs WHERE agent_id IN (%s) AND task_id IN (%s)`,
|
||
strings.Repeat("?,", len(agentIDs)-1)+"?",
|
||
strings.Repeat("?,", len(taskIDs)-1)+"?",
|
||
)
|
||
rows, err := d.Query(query, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
out := make(map[string]time.Time, len(agentIDs)*len(taskIDs))
|
||
for rows.Next() {
|
||
var agentID, taskID string
|
||
var ts time.Time
|
||
if err := rows.Scan(&agentID, &taskID, &ts); err != nil {
|
||
return nil, err
|
||
}
|
||
out[agentID+":"+taskID] = ts
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
func scanFleetTaskRow(row *sql.Row) (*models.FleetTask, error) {
|
||
t := &models.FleetTask{}
|
||
var enabled int
|
||
err := row.Scan(&t.ID, &t.Name, &enabled, &t.Trigger, &t.IntervalHours, &t.CronTime, &t.Action, &t.Command, &t.Target, &t.CreatedAt, &t.UpdatedAt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
t.Enabled = enabled == 1
|
||
return t, nil
|
||
}
|
||
|
||
func scanFleetTasks(rows *sql.Rows) ([]*models.FleetTask, error) {
|
||
var out []*models.FleetTask
|
||
for rows.Next() {
|
||
t := &models.FleetTask{}
|
||
var enabled int
|
||
if err := rows.Scan(&t.ID, &t.Name, &enabled, &t.Trigger, &t.IntervalHours, &t.CronTime, &t.Action, &t.Command, &t.Target, &t.CreatedAt, &t.UpdatedAt); err != nil {
|
||
return nil, err
|
||
}
|
||
t.Enabled = enabled == 1
|
||
out = append(out, t)
|
||
}
|
||
return out, rows.Err()
|
||
}
|