package monitor import ( "context" "encoding/json" "log" "sync" "time" "yaum/internal/alerts" "yaum/internal/model" "yaum/internal/sse" "yaum/internal/store" ) type checkResult struct { Heartbeat model.Heartbeat Service model.Service } type Worker struct { store store.Store interval time.Duration timeout time.Duration results chan checkResult broadcaster *sse.Broadcaster smtpCfg *alerts.SMTPConfig quit chan struct{} mu sync.Mutex lastChecked map[int]time.Time } func NewWorker(s store.Store, interval, timeout time.Duration, b *sse.Broadcaster) *Worker { return &Worker{ store: s, interval: interval, timeout: timeout, results: make(chan checkResult, 100), broadcaster: b, quit: make(chan struct{}), lastChecked: make(map[int]time.Time), } } func (w *Worker) WithSMTP(cfg *alerts.SMTPConfig) *Worker { w.smtpCfg = cfg return w } func (w *Worker) Start(ctx context.Context) { granularity := w.interval if granularity > 10*time.Second { granularity = 10 * time.Second } ticker := time.NewTicker(granularity) defer ticker.Stop() log.Printf("[worker] iniciado | global_interval=%v granularity=%v timeout=%v", w.interval, granularity, w.timeout) go w.resultConsumer() w.runChecks(ctx) for { select { case <-ticker.C: w.runChecks(ctx) case <-w.quit: log.Println("[worker] encerrando...") return } } } func (w *Worker) runChecks(ctx context.Context) { services, err := w.store.ListActiveServices() if err != nil { log.Printf("[worker] erro ao listar servicos ativos: %v", err) return } now := time.Now() var due []model.Service w.mu.Lock() for _, svc := range services { svcInterval := time.Duration(svc.IntervalSeconds) * time.Second if svcInterval <= 0 { svcInterval = w.interval } last, ok := w.lastChecked[svc.ID] if !ok || now.Sub(last) >= svcInterval { w.lastChecked[svc.ID] = now due = append(due, svc) } } w.mu.Unlock() if len(due) == 0 { return } log.Printf("[worker] verificando %d servico(s) de %d ativo(s)", len(due), len(services)) for _, svc := range due { go w.checkService(ctx, svc) } } func (w *Worker) resultConsumer() { for { select { case res := <-w.results: hb := res.Heartbeat svc := res.Service w.checkTransition(svc, hb) if err := w.store.AddHeartbeat(hb); err != nil { log.Printf("[worker] erro ao salvar heartbeat: %v", err) } if w.broadcaster != nil { data, err := json.Marshal(hb) if err == nil { w.broadcaster.Publish(data) } } status := "UP" if !hb.IsUp { status = "DOWN" } log.Printf("[worker] service=%d status=%s code=%d time=%dms", hb.ServiceID, status, hb.StatusCode, hb.ResponseTimeMs) case <-w.quit: return } } } func (w *Worker) checkTransition(svc model.Service, hb model.Heartbeat) { prev, err := w.store.GetLastHeartbeat(svc.ID) if err != nil || prev == nil { return } if prev.IsUp == hb.IsUp { return } log.Printf("[worker] transicao de status service=%d: %v -> %v", svc.ID, prev.IsUp, hb.IsUp) w.fireAlerts(svc, hb, prev.IsUp) } func (w *Worker) fireAlerts(svc model.Service, hb model.Heartbeat, wasUp bool) { if svc.DiscordWebhookURL != "" { go alerts.SendDiscord(svc.DiscordWebhookURL, svc, hb, wasUp) } if svc.AlertEmail != "" && w.smtpCfg != nil { go alerts.SendEmail(svc, hb, wasUp, w.smtpCfg) } } func (w *Worker) checkService(ctx context.Context, svc model.Service) { checkCtx, cancel := context.WithTimeout(ctx, w.timeout) defer cancel() hb, ssl := CheckHTTP(checkCtx, svc.URL, svc.KeywordToFind) hb.ServiceID = svc.ID if ssl != nil { ssl.ServiceID = svc.ID if err := w.store.SaveSSLInfo(*ssl); err != nil { log.Printf("[worker] erro ao salvar ssl info: %v", err) } } w.results <- checkResult{Heartbeat: hb, Service: svc} } func (w *Worker) Stop() { close(w.quit) }