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() 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 { last, ok := s.db.LastFleetTaskRun(agentID, t.ID) if 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) }