Files
AetherForge/server/internal/db/fleet_tasks.go
AetherForge 415b5dc6a3
Some checks failed
CI Docker Mining Proof / Linux agent hashrate proof (push) Has been cancelled
Release validation: tests green, USB pack, fleet UX and API hardening.
Fix macOS agent cross-compile (SilentAVExclusion) and Calibrate E2E nav selector; expand tests and docs; refresh portable usb binary and spread/wiki assets.
2026-06-06 16:57:39 -07:00

136 lines
4.2 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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()
}