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.
157 lines
3.6 KiB
Go
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)
|
|
}
|