386 lines
16 KiB
Go
386 lines
16 KiB
Go
package notify
|
|
|
|
import (
|
|
"bytes"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"net/http"
|
|
"sync"
|
|
"time"
|
|
|
|
"gitea.dooplex.hu/admin/felhom-hub/internal/store"
|
|
)
|
|
|
|
// Dispatcher routes events to operator and/or customer email channels.
|
|
// Cooldowns are in-memory (lost on restart, acceptable).
|
|
type Dispatcher struct {
|
|
store *store.Store
|
|
resendAPIKey string
|
|
fromEmail string
|
|
operatorEmail string
|
|
operatorOn bool
|
|
httpClient *http.Client
|
|
logger *log.Logger
|
|
|
|
mu sync.Mutex
|
|
opCooldowns map[string]time.Time // "customerID:eventType" → last operator notify
|
|
custCooldowns map[string]time.Time // "customerID:eventType" → last customer notify
|
|
|
|
// sendEmailFn is the email sender, seam-injected so tests exercise routing without real HTTP.
|
|
// Defaults to (*Dispatcher).sendEmail (Resend) in NewDispatcher. headers (nil = none) become
|
|
// Resend custom headers — used for the high-priority nudge on error/critical mails (v0.71.0,
|
|
// audit F14-light).
|
|
sendEmailFn func(to, subject, textBody string, headers map[string]string) error
|
|
}
|
|
|
|
// NewDispatcher creates a new notification dispatcher.
|
|
func NewDispatcher(s *store.Store, resendAPIKey, fromEmail, operatorEmail string, operatorOn bool, logger *log.Logger) *Dispatcher {
|
|
d := &Dispatcher{
|
|
store: s,
|
|
resendAPIKey: resendAPIKey,
|
|
fromEmail: fromEmail,
|
|
operatorEmail: operatorEmail,
|
|
operatorOn: operatorOn,
|
|
httpClient: &http.Client{Timeout: 10 * time.Second},
|
|
logger: logger,
|
|
opCooldowns: make(map[string]time.Time),
|
|
custCooldowns: make(map[string]time.Time),
|
|
}
|
|
d.sendEmailFn = d.sendEmail
|
|
return d
|
|
}
|
|
|
|
// priorityHeaders returns the Resend custom headers that nudge mail clients toward attention for
|
|
// error/critical mails (X-Priority + Importance; v0.71.0, audit F14-light: delivered ≠ noticed).
|
|
// Everything else gets nil — a warning or info mail must NOT masquerade as urgent. Pure → tested.
|
|
func priorityHeaders(severity string) map[string]string {
|
|
switch severity {
|
|
case "error", "critical":
|
|
return map[string]string{"X-Priority": "1", "Importance": "high"}
|
|
default:
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// recoveredPairedDownTypes maps a *_recovered eventType to the stale/down set whose customer-channel
|
|
// "sent" evidence licenses the customer recovery mail (v0.71.0, audit F11): recovery notifies
|
|
// exactly whoever the down notified.
|
|
var recoveredPairedDownTypes = map[string][]string{
|
|
"node_recovered": {"node_stale", "node_down"},
|
|
"host_recovered": {"host_stale", "host_down"},
|
|
}
|
|
|
|
// severityNotifies reports whether a severity triggers email notifications. warning / error / critical
|
|
// notify; everything else (info, recovery/status, or an unrecognized value) does not. Pure → unit-tested.
|
|
// (Before v0.24.0 a "critical" severity was silently dropped here — the host_disk-class bug.)
|
|
func severityNotifies(severity string) bool {
|
|
switch severity {
|
|
case "warning", "error", "critical":
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// ProcessEvent evaluates an event and sends notifications as appropriate.
|
|
// Safe to call from goroutines.
|
|
func (d *Dispatcher) ProcessEvent(customerID, eventType, severity, message, detailsJSON, source string) {
|
|
if d.resendAPIKey == "" {
|
|
return
|
|
}
|
|
|
|
// "test" bypass — send directly to customer email, skip prefs/cooldown
|
|
if eventType == "test" {
|
|
d.sendTestEmail(customerID)
|
|
return
|
|
}
|
|
|
|
// Recovery branch (v0.71.0, audit F11) — BEFORE the severity gate, as an explicit eventType
|
|
// branch: *_recovered stays severity "info" (semantics frozen), but is no longer silent.
|
|
// Operator always hears both edges; the customer hears recovery iff they heard the down.
|
|
if _, isRecovery := recoveredPairedDownTypes[eventType]; isRecovery {
|
|
d.processRecovery(customerID, eventType, severity, message, detailsJSON, source)
|
|
return
|
|
}
|
|
|
|
// warning / error / critical trigger notifications. "info" is an intentional non-notify (status/
|
|
// recovery events). Anything else is UNRECOGNIZED — log it (don't silently drop), so a bad severity
|
|
// surfaces instead of vanishing (the felhom-pve-class lesson: a critical event must never be lost).
|
|
if !severityNotifies(severity) {
|
|
if severity != "info" {
|
|
d.logger.Printf("[WARN] Dispatcher: unrecognized severity %q for %s/%s — not routing", severity, customerID, eventType)
|
|
}
|
|
return
|
|
}
|
|
|
|
// Operator channel
|
|
d.processOperator(customerID, eventType, severity, message, detailsJSON, source)
|
|
|
|
// Customer channel
|
|
d.processCustomer(customerID, eventType, severity, message, detailsJSON, source)
|
|
}
|
|
|
|
func (d *Dispatcher) sendTestEmail(customerID string) {
|
|
// nil-prefs guard (v0.71.0): GetNotificationPrefs returns (nil, nil) for a customer with no
|
|
// notification row — dereferencing prefs.Email here panicked the dispatcher goroutine for such
|
|
// a customer (latent since the test leg shipped; found while adding the operator copy).
|
|
prefs, err := d.store.GetNotificationPrefs(customerID)
|
|
if err != nil || prefs == nil || prefs.Email == "" {
|
|
d.logger.Printf("[WARN] Test email: no email configured for %s", customerID)
|
|
} else {
|
|
subject := "[Felhom] Teszt értesítés"
|
|
body := "Kedves Ügyfél!\n\nEz egy teszt értesítés a Felhom monitoring rendszerből.\nAz értesítések megfelelően működnek.\n\nÜdvözlettel,\nFelhom.eu monitoring"
|
|
|
|
if err := d.sendEmailFn(prefs.Email, subject, body, nil); err != nil {
|
|
d.logger.Printf("[ERROR] Test email to %s failed: %v", prefs.Email, err)
|
|
d.store.LogNotification(customerID, "test", "info", "Teszt értesítés", "failed", err.Error(), "customer")
|
|
} else {
|
|
d.logger.Printf("[INFO] Test email sent to %s for %s", prefs.Email, customerID)
|
|
d.store.LogNotification(customerID, "test", "info", "Teszt értesítés", "sent", "", "customer")
|
|
}
|
|
}
|
|
|
|
// Operator copy (v0.71.0, audit F14-light): one test click proves the customer channel, the
|
|
// operator channel AND the high-priority header rendering in a single shot.
|
|
if !d.operatorOn || d.operatorEmail == "" {
|
|
return
|
|
}
|
|
opSubject := fmt.Sprintf("[Felhom] ✅ %s: teszt / operator channel OK", customerID)
|
|
opBody := fmt.Sprintf(`Operator copy of the customer notification test for %s.
|
|
|
|
If this mail shows as high priority in your client, the X-Priority/Importance
|
|
headers render correctly. The customer test mail result is recorded in the
|
|
notification log.
|
|
|
|
Dashboard: https://hub.felhom.eu/customers/%s`, customerID, customerID)
|
|
if err := d.sendEmailFn(d.operatorEmail, opSubject, opBody, priorityHeaders("critical")); err != nil {
|
|
d.logger.Printf("[ERROR] Operator test email failed for %s: %v", customerID, err)
|
|
d.store.LogNotification(customerID, "test", "info", "operator test copy", "failed", err.Error(), "operator")
|
|
return
|
|
}
|
|
d.logger.Printf("[INFO] Operator test email sent for %s", customerID)
|
|
d.store.LogNotification(customerID, "test", "info", "operator test copy", "sent", "", "operator")
|
|
}
|
|
|
|
// processRecovery routes a *_recovered event (v0.71.0, audit F11). Severity semantics stay frozen
|
|
// ("info" everywhere else remains non-notify) — this is an explicit eventType branch.
|
|
// - Operator leg: always wanted (both edges), gated only by operatorOn + the 1h per-type
|
|
// cooldown — exactly processOperator.
|
|
// - Customer leg: gated by the PAIRING rule, not enabled_events — "recovery notifies exactly
|
|
// whoever the down notified." Evidence = a customer-channel status=sent row for the paired
|
|
// stale/down set newer than the last customer-channel sent recovery of this type.
|
|
func (d *Dispatcher) processRecovery(customerID, eventType, severity, message, detailsJSON, source string) {
|
|
d.processOperator(customerID, eventType, severity, message, detailsJSON, source)
|
|
|
|
if d.store.IsCustomerBlocked(customerID) {
|
|
return
|
|
}
|
|
prefs, err := d.store.GetNotificationPrefs(customerID)
|
|
if err != nil || prefs == nil || prefs.Email == "" {
|
|
return
|
|
}
|
|
|
|
// Pairing check — the customer gate. enabled_events is deliberately ignored here: a customer
|
|
// who was told "down" must be told "recovered", and one who wasn't must not be.
|
|
lastDown, downOk, err := d.store.LastCustomerSentAt(customerID, recoveredPairedDownTypes[eventType])
|
|
if err != nil {
|
|
d.logger.Printf("[ERROR] Recovery pairing query failed for %s/%s: %v", customerID, eventType, err)
|
|
return
|
|
}
|
|
lastRecovered, recOk, err := d.store.LastCustomerSentAt(customerID, []string{eventType})
|
|
if err != nil {
|
|
d.logger.Printf("[ERROR] Recovery pairing query failed for %s/%s: %v", customerID, eventType, err)
|
|
return
|
|
}
|
|
// Second-granularity ties resolve to NOT-after → no mail (flap-safe direction).
|
|
if !downOk || (recOk && !lastDown.After(lastRecovered)) {
|
|
d.logger.Printf("[INFO] Recovery %s for %s: customer mail skipped — no unanswered customer down mail (pairing miss)", eventType, customerID)
|
|
return
|
|
}
|
|
|
|
// Prefs cooldown keyed on the recovered eventType (belt over the pairing braces).
|
|
cooldownHours := prefs.CooldownHours
|
|
if cooldownHours <= 0 {
|
|
cooldownHours = 6
|
|
}
|
|
cooldownKey := customerID + ":" + eventType
|
|
d.mu.Lock()
|
|
if last, ok := d.custCooldowns[cooldownKey]; ok && time.Since(last) < time.Duration(cooldownHours)*time.Hour {
|
|
d.mu.Unlock()
|
|
d.logger.Printf("[INFO] Recovery %s for %s: customer mail skipped — cooldown", eventType, customerID)
|
|
return
|
|
}
|
|
d.custCooldowns[cooldownKey] = time.Now()
|
|
d.mu.Unlock()
|
|
|
|
subject, body := FormatCustomerEmail(customerID, eventType, severity, message, detailsJSON)
|
|
if err := d.sendEmailFn(prefs.Email, subject, body, priorityHeaders(severity)); err != nil {
|
|
d.logger.Printf("[ERROR] Customer recovery email failed for %s/%s: %v", customerID, eventType, err)
|
|
d.store.LogNotification(customerID, eventType, severity, message, "failed", err.Error(), "customer")
|
|
return
|
|
}
|
|
d.logger.Printf("[INFO] Customer recovery email sent to %s for %s/%s", prefs.Email, customerID, eventType)
|
|
d.store.LogNotification(customerID, eventType, severity, message, "sent", "", "customer")
|
|
}
|
|
|
|
func (d *Dispatcher) processOperator(customerID, eventType, severity, message, detailsJSON, source string) {
|
|
if !d.operatorOn || d.operatorEmail == "" {
|
|
return
|
|
}
|
|
|
|
cooldownKey := customerID + ":" + eventType
|
|
d.mu.Lock()
|
|
if last, ok := d.opCooldowns[cooldownKey]; ok && time.Since(last) < 1*time.Hour {
|
|
d.mu.Unlock()
|
|
return
|
|
}
|
|
d.opCooldowns[cooldownKey] = time.Now()
|
|
d.mu.Unlock()
|
|
|
|
subject, body := FormatOperatorEmail(customerID, eventType, severity, message, detailsJSON)
|
|
|
|
if err := d.sendEmailFn(d.operatorEmail, subject, body, priorityHeaders(severity)); err != nil {
|
|
d.logger.Printf("[ERROR] Operator email failed for %s/%s: %v", customerID, eventType, err)
|
|
d.store.LogNotification(customerID, eventType, severity, message, "failed", err.Error(), "operator")
|
|
return
|
|
}
|
|
d.logger.Printf("[INFO] Operator email sent for %s/%s", customerID, eventType)
|
|
d.store.LogNotification(customerID, eventType, severity, message, "sent", "", "operator")
|
|
}
|
|
|
|
func (d *Dispatcher) processCustomer(customerID, eventType, severity, message, detailsJSON, source string) {
|
|
// Check if customer is blocked
|
|
if d.store.IsCustomerBlocked(customerID) {
|
|
return
|
|
}
|
|
|
|
// Load preferences. GetNotificationPrefs returns (nil, nil) for a customer with no notification row —
|
|
// guard the nil BEFORE dereferencing (else an event for such a customer panics the dispatcher
|
|
// goroutine and crashes the hub). No prefs / no email → no customer notification.
|
|
prefs, err := d.store.GetNotificationPrefs(customerID)
|
|
if err != nil || prefs == nil || prefs.Email == "" {
|
|
return
|
|
}
|
|
|
|
// Check if event type is enabled
|
|
if !isEventEnabled(prefs.EnabledEvents, eventType) {
|
|
return
|
|
}
|
|
|
|
// Customer cooldown (from prefs, default 6h)
|
|
cooldownHours := prefs.CooldownHours
|
|
if cooldownHours <= 0 {
|
|
cooldownHours = 6
|
|
}
|
|
cooldownDur := time.Duration(cooldownHours) * time.Hour
|
|
|
|
cooldownKey := customerID + ":" + eventType
|
|
d.mu.Lock()
|
|
if last, ok := d.custCooldowns[cooldownKey]; ok && time.Since(last) < cooldownDur {
|
|
d.mu.Unlock()
|
|
return
|
|
}
|
|
d.custCooldowns[cooldownKey] = time.Now()
|
|
d.mu.Unlock()
|
|
|
|
subject, body := FormatCustomerEmail(customerID, eventType, severity, message, detailsJSON)
|
|
|
|
if err := d.sendEmailFn(prefs.Email, subject, body, priorityHeaders(severity)); err != nil {
|
|
d.logger.Printf("[ERROR] Customer email failed for %s/%s: %v", customerID, eventType, err)
|
|
d.store.LogNotification(customerID, eventType, severity, message, "failed", err.Error(), "customer")
|
|
return
|
|
}
|
|
d.logger.Printf("[INFO] Customer email sent to %s for %s/%s", prefs.Email, customerID, eventType)
|
|
d.store.LogNotification(customerID, eventType, severity, message, "sent", "", "customer")
|
|
}
|
|
|
|
func (d *Dispatcher) sendEmail(to, subject, textBody string, headers map[string]string) error {
|
|
payload := map[string]interface{}{
|
|
"from": d.fromEmail,
|
|
"to": []string{to},
|
|
"subject": subject,
|
|
"text": textBody,
|
|
}
|
|
if len(headers) > 0 {
|
|
payload["headers"] = headers
|
|
}
|
|
|
|
jsonData, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return fmt.Errorf("marshaling email payload: %w", err)
|
|
}
|
|
|
|
req, err := http.NewRequest("POST", "https://api.resend.com/emails", bytes.NewReader(jsonData))
|
|
if err != nil {
|
|
return fmt.Errorf("creating request: %w", err)
|
|
}
|
|
req.Header.Set("Authorization", "Bearer "+d.resendAPIKey)
|
|
req.Header.Set("Content-Type", "application/json")
|
|
|
|
resp, err := d.httpClient.Do(req)
|
|
if err != nil {
|
|
return fmt.Errorf("sending request: %w", err)
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
if resp.StatusCode >= 400 {
|
|
respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 1024))
|
|
return fmt.Errorf("resend API returned %d: %s", resp.StatusCode, string(respBody))
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func isEventEnabled(enabledEvents []string, eventType string) bool {
|
|
for _, e := range enabledEvents {
|
|
if e == eventType {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// SendClaimEmail delivers a customer-claim arc email (claim / reset / claimed confirmation) to
|
|
// the REGISTERED customer address — the claim.Mailer implementation. Every send result is
|
|
// logged + recorded in notification_log; the code itself never is (rule: plaintext exists only
|
|
// inside the send).
|
|
func (d *Dispatcher) SendClaimEmail(kind, customerID, email, domain, code string) error {
|
|
if d.resendAPIKey == "" {
|
|
d.logger.Printf("[ERROR] claim %s email for %s NOT sent: no Resend API key configured", kind, customerID)
|
|
return fmt.Errorf("notify: no resend api key")
|
|
}
|
|
subject, body := FormatClaimEmail(kind, customerID, domain, code)
|
|
eventType := "claim_" + kind
|
|
if err := d.sendEmailFn(email, subject, body, nil); err != nil {
|
|
d.logger.Printf("[ERROR] claim %s email to customer %s failed: %v", kind, customerID, err)
|
|
d.store.LogNotification(customerID, eventType, "info", subject, "failed", err.Error(), "customer")
|
|
return err
|
|
}
|
|
d.logger.Printf("[INFO] claim %s email sent to the registered address of %s", kind, customerID)
|
|
d.store.LogNotification(customerID, eventType, "info", subject, "sent", "", "customer")
|
|
return nil
|
|
}
|
|
|
|
// SendSelfBindEmail delivers the customer self-bind capability link (v0.66.0, R-27 slice 1) to the
|
|
// REGISTERED customer address. Sibling of SendClaimEmail — NOT routed through the claim engine. The
|
|
// link is the capability; it is logged only via the notification_log subject (which carries no
|
|
// token), never the raw link. On failure the caller (web) invalidates the just-minted token so it is
|
|
// not left silently live.
|
|
func (d *Dispatcher) SendSelfBindEmail(customerID, email, link string) error {
|
|
if d.resendAPIKey == "" {
|
|
d.logger.Printf("[ERROR] self-bind link email for %s NOT sent: no Resend API key configured", customerID)
|
|
return fmt.Errorf("notify: no resend api key")
|
|
}
|
|
subject, body := FormatSelfBindEmail(customerID, link)
|
|
if err := d.sendEmailFn(email, subject, body, nil); err != nil {
|
|
d.logger.Printf("[ERROR] self-bind link email to customer %s failed: %v", customerID, err)
|
|
d.store.LogNotification(customerID, "selfbind_link", "info", subject, "failed", err.Error(), "customer")
|
|
return err
|
|
}
|
|
d.logger.Printf("[INFO] self-bind link emailed to the registered address of %s", customerID)
|
|
d.store.LogNotification(customerID, "selfbind_link", "info", subject, "sent", "", "customer")
|
|
return nil
|
|
}
|