d3b284863e
gates / gates (push) Successful in 26s
Allowlisted, not operator-only, seeded for new households and added once (add-only) to every existing enabled_events row. mail.event entries in hu and en name the app from details.stack_name. Per-app cooldown on both the operator and the household leg. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
809 lines
41 KiB
Go
809 lines
41 KiB
Go
package notify
|
|
|
|
import (
|
|
"bytes"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"net/http"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"gitea.dooplex.hu/admin/felhom-hub/internal/store"
|
|
|
|
"gitea.dooplex.hu/admin/felhom-hub/internal/i18n"
|
|
)
|
|
|
|
// 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"},
|
|
// R-97a. This branch runs BEFORE the severity gate, which is exactly why the recovery belongs
|
|
// here: `whole_guest_backup_recovered` is severity "info", and severityNotifies drops "info", so
|
|
// routing it normally would store the event and never mail it — the operator would be told the
|
|
// tier broke and never told it healed, which is the half of Scenario B that matters.
|
|
//
|
|
// The customer leg needs no special handling: it is PAIRING-gated on a customer-channel "sent"
|
|
// row for the down type, and `whole_guest_backup_failed` has no customerMessages entry and is in
|
|
// nobody's enabled_events — so no such row can exist, and the customer correctly hears neither
|
|
// edge. Operator hears both.
|
|
"whole_guest_backup_recovered": {"whole_guest_backup_failed"},
|
|
// R-339 (v0.106.0) — box REACHABILITY all-clears. Same reasoning as the entry above and the same
|
|
// necessity: both are severity "info", so without an entry here severityNotifies drops them and
|
|
// the operator is told the off-site tier went blind and never told it came back.
|
|
//
|
|
// The customer leg is a no-op BY CONSTRUCTION, not by luck: both scopes are customer-less
|
|
// ("pbsdr-box" / "pool-box" are not customer IDs), so GetNotificationPrefs finds no row and
|
|
// LastCustomerSentAt can never locate a paired down row for them. Operator hears both edges;
|
|
// no customer hears either, which is correct — neither store belongs to a customer.
|
|
"pbsdr_box_recovered": {"pbsdr_box_unreachable"},
|
|
"offsite_box_recovered": {"offsite_box_unreachable"},
|
|
}
|
|
|
|
// 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.
|
|
//
|
|
// This is the HUB-GENERATED entry point: the hub composed `message` itself (a checker, a staleness
|
|
// sweep, an escrow decision), so there is no second sentence to carry and the customer mail is
|
|
// rendered from the bundle alone.
|
|
func (d *Dispatcher) ProcessEvent(customerID, eventType, severity, message, detailsJSON, source string) {
|
|
d.processEvent(customerID, eventType, severity, message, "", detailsJSON, source)
|
|
}
|
|
|
|
// ProcessBoxEvent is the entry point for an event a BOX sent (POST /api/v1/event).
|
|
//
|
|
// `messageCustomer` is the box's own sentence in the household's language, sent beside the Hungarian
|
|
// `message` by controller v0.256.0 and later (R-558). It is optional forever — an older box sends
|
|
// none and the mail is then exactly what it was.
|
|
//
|
|
// Kept SEPARATE from ProcessEvent rather than added as a parameter to it: 44 of the 45 call sites
|
|
// are hub-generated and have no such sentence, and widening all of them would have meant 44 edits
|
|
// whose only content is an empty string — churn that hides the one call site that matters.
|
|
func (d *Dispatcher) ProcessBoxEvent(customerID, eventType, severity, message, messageCustomer, detailsJSON, source string) {
|
|
d.processEvent(customerID, eventType, severity, message, messageCustomer, detailsJSON, source)
|
|
}
|
|
|
|
func (d *Dispatcher) processEvent(customerID, eventType, severity, message, messageCustomer, 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, messageCustomer, detailsJSON, source)
|
|
return
|
|
}
|
|
|
|
// R-182: record-only types are written down and never mailed. Placed BEFORE the severity gate
|
|
// so the row is written whatever the severity — the record must not inherit the notification's
|
|
// conditions, which is the coupling this whole finding is about.
|
|
if recordOnlyEvents[eventType] {
|
|
if err := d.store.LogNotification(customerID, eventType, severity, message, "recorded",
|
|
"record-only: the per-run digest (backup_run_failures) carries the notification", "operator"); err != nil {
|
|
d.logger.Printf("[WARN] Failed to record %s for %s: %v", eventType, customerID, err)
|
|
}
|
|
d.logger.Printf("[INFO] Recorded (not mailed) %s for %s — the run digest is the notification", eventType, customerID)
|
|
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).
|
|
//
|
|
// R-387 — THIS BRANCH IS **KEPT DELIBERATELY**, and here is why, because the question was asked
|
|
// and a branch that cannot execute without a note is the thing to avoid.
|
|
//
|
|
// For an event arriving over the API it is genuinely unreachable: the ingest handler coerces any
|
|
// unknown severity to "info" before this is called, so the one value it looks for cannot arrive.
|
|
// **But the API is not the only producer.** `cmd/hub/main.go` wires `dispatcher.ProcessEvent`
|
|
// DIRECTLY as the `monitor.EventNotifyFunc` for the staleness, host-staleness and offsite-box
|
|
// checkers, and those hub-generated events never pass through the handler at all. For every one
|
|
// of them this line is the ONLY severity guard there is.
|
|
//
|
|
// Removing it as "dead" would therefore have deleted the live half while leaving the dead half
|
|
// looking like the reason. Verified 2026-08-23: every severity literal in `internal/monitor` (90
|
|
// of them) is already in the vocabulary — so the guard is currently silent because the producers
|
|
// are correct, which is exactly what a guard looks like when it is working.
|
|
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, messageCustomer, 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 {
|
|
// The test mail follows the household's language like every other customer mail (R-558).
|
|
//
|
|
// It was the ONE that would not have: it is the only customer mail with its own hardcoded
|
|
// text, because it never goes through FormatCustomerEmail — and it is also the only one an
|
|
// operator can trigger on demand, so it is the one most likely to be used to CHECK whether
|
|
// the localisation works. Left alone, pressing "send test" for an English household would
|
|
// have answered that question wrongly.
|
|
b := i18n.Shared()
|
|
lang := d.store.CustomerLanguage(customerID)
|
|
subject := b.Msg(lang, "mail.test.subject")
|
|
body := b.Msg(lang, "mail.test.body")
|
|
|
|
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, messageCustomer, 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(d.store.CustomerLanguage(customerID),
|
|
customerID, eventType, severity, message, messageCustomer, 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")
|
|
}
|
|
|
|
// cooldownTierSuffix returns ":"+tier when the event's details carry a non-empty `tier`, else "".
|
|
//
|
|
// R-97a. The operator cooldown was keyed `customerID + ":" + eventType` alone, which is correct for
|
|
// every event that describes ONE thing — but a whole-guest backup failure describes ONE TIER, and a
|
|
// box has two. `felhom-pbs` failing at 09:00 would swallow `local` failing at 09:20 for the whole
|
|
// hour, so the operator would be told about the offsite tier and never about the local one. That is
|
|
// precisely the masking the per-tier signal exists to prevent.
|
|
//
|
|
// NARROW ON PURPOSE: the suffix is empty unless the producer opts in by sending a `tier`, so no
|
|
// existing event type's cooldown behaviour changes. Widening the key for everything would, e.g.,
|
|
// turn one hourly `app_start_failed` into one per app — a flood, not a fix.
|
|
func cooldownTierSuffix(detailsJSON string) string {
|
|
if detailsJSON == "" || !strings.Contains(detailsJSON, "\"tier\"") {
|
|
return ""
|
|
}
|
|
var d struct {
|
|
Tier string `json:"tier"`
|
|
}
|
|
if err := json.Unmarshal([]byte(detailsJSON), &d); err != nil || d.Tier == "" {
|
|
return ""
|
|
}
|
|
return ":" + d.Tier
|
|
}
|
|
|
|
// cooldownRunSuffix returns ":"+run_id when the event's details carry a non-empty `run_id`, else "".
|
|
//
|
|
// R-182. `cooldownTierSuffix`'s sibling, and deliberately a SEPARATE function rather than an extra
|
|
// branch inside it: `tier` keeps byte-identical semantics for every type that uses it, so R-97a's
|
|
// behaviour and its tests are untouched by this.
|
|
//
|
|
// WHY A BACKUP RUN NEEDS ONE. The run digest describes ONE RUN, and a box can have two in a day —
|
|
// the nightly one and a manual one the operator triggered *because* something looked wrong. With no
|
|
// run-scoped discriminator the 1-hour cooldown would swallow the second, which is the failure this
|
|
// row exists to fix, reappearing one level up: the operator presses the button, the run fails, and
|
|
// they are told nothing because the machine already wrote that hour.
|
|
//
|
|
// IT MAKES THE COOLDOWN EFFECTIVELY INERT FOR THIS TYPE, AND THAT IS THE INTENT, NOT AN OVERSIGHT.
|
|
// A digest is already rate-limited by construction — one per run, emitted only when something
|
|
// failed — so there is nothing for a timer to collapse. The cooldown protects against a repeating
|
|
// identical alert; a digest cannot repeat, because each run is a different run.
|
|
//
|
|
// NARROW, LIKE ITS SIBLING: empty unless the producer opts in by sending a `run_id`, so no existing
|
|
// event type's cooldown behaviour changes.
|
|
func cooldownRunSuffix(detailsJSON string) string {
|
|
if detailsJSON == "" || !strings.Contains(detailsJSON, "\"run_id\"") {
|
|
return ""
|
|
}
|
|
var d struct {
|
|
RunID string `json:"run_id"`
|
|
}
|
|
if err := json.Unmarshal([]byte(detailsJSON), &d); err != nil || d.RunID == "" {
|
|
return ""
|
|
}
|
|
return ":" + d.RunID
|
|
}
|
|
|
|
// perAppCooldownEvents is the register of event types whose operator cooldown is keyed PER APP.
|
|
//
|
|
// R-389. **This is a named allow-list and not a behaviour inferred from the payload, deliberately.**
|
|
// Several event types carry `stack_name` and must NOT be split per app — see the fence below — so a
|
|
// rule of the form "if it has a stack_name, split it" would silently change them. The register makes
|
|
// the decision reviewable one line at a time, exactly as `operatorOnlyEvents` does.
|
|
//
|
|
// ── THE FENCE, RECORDED SO IT CAN BE NARROWED LATER IF IT IS EVER WRONG ──────────────────────
|
|
//
|
|
// The backup family's cooldown is coarse **on purpose**. R-97a and R-182 exist precisely so that one
|
|
// full disk produces ONE mail listing every affected app, rather than one mail per app. Adding
|
|
// `stack_name` to the key for those types would undo both, and it would do it silently — the code
|
|
// would look more precise while the operator's inbox got twenty times louder.
|
|
//
|
|
// It is not hypothetical: **`crossdrive_failed` is severity `error`, reaches the operator leg, and
|
|
// carries `stack_name`** through a different struct (`CrossDriveDetails`, not `AppDetails`). A
|
|
// payload-shape rule would have caught it and split it. This register does not.
|
|
//
|
|
// The fenced ACT is *adding an entry here for a type whose family has a digest or a coarse-by-design
|
|
// cooldown*. Adding one for a type that genuinely has no digest and alarms per app is the intended
|
|
// use.
|
|
var perAppCooldownEvents = map[string]bool{
|
|
// The ONLY member as of hub v0.108.0. An app going down is a per-app fault with no digest: there
|
|
// is no `apps_down_run` summarising a scan the way `backup_run_failures` summarises a run, so
|
|
// per-app is the only grain available that does not lose alarms. Measured 2026-08-23: two apps
|
|
// four minutes apart produced one mail and one suppression.
|
|
"app_start_failed": true,
|
|
// v0.120.0: an update outcome is per app with no digest — two apps on one night are two alarms.
|
|
"app_update_undone": true,
|
|
"app_update_held": true,
|
|
}
|
|
|
|
// perAppCustomerCooldownEvents is the CUSTOMER-leg sibling of perAppCooldownEvents (v0.120.0). The
|
|
// household's cooldown was keyed `customer:type` for every type, so two apps undone on the same night
|
|
// would reach the household as ONE mail. Its own register, deliberately — putting app_start_failed's
|
|
// customer leg on a per-app key would change a type nobody asked to change.
|
|
var perAppCustomerCooldownEvents = map[string]bool{
|
|
"app_update_undone": true,
|
|
"app_update_held": true,
|
|
}
|
|
|
|
// cooldownStackSuffix returns ":"+stack_name when the event's details carry a non-empty `stack_name`
|
|
// AND the event type is in `perAppCooldownEvents`, else "".
|
|
//
|
|
// R-389. The third sibling of `cooldownTierSuffix` and `cooldownRunSuffix`, and deliberately a
|
|
// SEPARATE function for the reason `cooldownRunSuffix`'s docstring already gives: the existing two
|
|
// keep byte-identical semantics for every type that uses them, so R-97a's and R-182's behaviour and
|
|
// their tests are untouched by this.
|
|
//
|
|
// WHY AN APP NEEDS ONE. The operator cooldown collapses everything sharing a key for an hour. Until
|
|
// v0.108.0 the key named the event TYPE and not the app, so a second app going down inside that hour
|
|
// was recorded `suppressed` and never mailed. Measured live 2026-08-23: `bookstack` sent at 09:27:51,
|
|
// `privatebin` suppressed at 09:31:51 under `key=demo-hp:app_start_failed`. **Three apps dying
|
|
// together produced one mail.**
|
|
//
|
|
// IT TAKES THE EVENT TYPE AS WELL AS THE DETAILS, unlike its two siblings, and that asymmetry is the
|
|
// whole safety property — see `perAppCooldownEvents`. The siblings can be payload-driven because
|
|
// `tier` and `run_id` appear only on types that want that grain; `stack_name` does not have that
|
|
// property.
|
|
//
|
|
// NARROW AND FAIL-SOFT, like its siblings: empty on an absent, malformed or empty value, so a
|
|
// degraded payload falls back to today's key and **the mail still goes**. Losing an alarm is worse
|
|
// than mis-routing one.
|
|
func cooldownStackSuffix(eventType, detailsJSON string) string {
|
|
if !perAppCooldownEvents[eventType] {
|
|
return ""
|
|
}
|
|
return cooldownStackSuffixFor(detailsJSON)
|
|
}
|
|
|
|
// cooldownStackSuffixFor is the payload half of cooldownStackSuffix, register-free — callers decide
|
|
// which register applies (operator: perAppCooldownEvents; customer: perAppCustomerCooldownEvents).
|
|
func cooldownStackSuffixFor(detailsJSON string) string {
|
|
if detailsJSON == "" || !strings.Contains(detailsJSON, "\"stack_name\"") {
|
|
return ""
|
|
}
|
|
var d struct {
|
|
StackName string `json:"stack_name"`
|
|
}
|
|
if err := json.Unmarshal([]byte(detailsJSON), &d); err != nil || d.StackName == "" {
|
|
return ""
|
|
}
|
|
return ":" + d.StackName
|
|
}
|
|
|
|
// nodeLivenessEvents skip the 1-hour operator cooldown (OPERATOR RULING 2026-09-15, decision A;
|
|
// 08-alarm-ladder.md §5). BIGNIGHT F9: the box was dead for 33 minutes and the `node_stale` mail was
|
|
// suppressed because F8's `node_stale` had used the hour 39 minutes earlier; the `node_recovered`
|
|
// mail was suppressed the same way. "The box is down" must not wait out a quiet hour. A 5-minute
|
|
// dedupe stays, so a flapping link cannot mail every sweep. The key is unchanged (customer:type), and
|
|
// a customer has one box, so the dedupe is per host. EXTENDED 2026-09-16 (ruling 2, R-529) to the
|
|
// agent-plane siblings host_stale / host_down / host_recovered — same box, same sentence.
|
|
var nodeLivenessEvents = map[string]bool{
|
|
"node_stale": true,
|
|
"node_down": true,
|
|
"node_recovered": true,
|
|
// OPERATOR RULING 2026-09-16 (ruling 2, hub v0.115.0): the HOST-plane siblings join it. They are
|
|
// the same sentence about the same box — the agent's dead-man's-switch rather than the
|
|
// controller's — and R-529 was filed precisely because the 2026-09-15 ruling named only node_*.
|
|
"host_stale": true,
|
|
"host_down": true,
|
|
"host_recovered": true,
|
|
}
|
|
|
|
const (
|
|
operatorCooldown = 1 * time.Hour
|
|
nodeLivenessDedupeWindow = 5 * time.Minute
|
|
)
|
|
|
|
// operatorCooldownFor returns the operator-mail cooldown for an event type. Pinned by
|
|
// TestOperatorCooldown_NodeLivenessBypassesQuietHour.
|
|
func operatorCooldownFor(eventType string) time.Duration {
|
|
if nodeLivenessEvents[eventType] {
|
|
return nodeLivenessDedupeWindow
|
|
}
|
|
return operatorCooldown
|
|
}
|
|
|
|
func (d *Dispatcher) processOperator(customerID, eventType, severity, message, detailsJSON, source string) {
|
|
if !d.operatorOn || d.operatorEmail == "" {
|
|
return
|
|
}
|
|
|
|
// R-389 added the third suffix. It is allow-listed to one event type, so every other type's key
|
|
// is byte-identical to v0.107.0's — pinned by TestR389_NoOtherEventTypeKeyChanges.
|
|
cooldownKey := customerID + ":" + eventType +
|
|
cooldownTierSuffix(detailsJSON) + cooldownRunSuffix(detailsJSON) +
|
|
cooldownStackSuffix(eventType, detailsJSON)
|
|
window := operatorCooldownFor(eventType)
|
|
d.mu.Lock()
|
|
if last, ok := d.opCooldowns[cooldownKey]; ok && time.Since(last) < window {
|
|
d.mu.Unlock()
|
|
// R-182: RECORD THE SUPPRESSION. This used to be a bare `return` — the event was dropped
|
|
// before any LogNotification, so a cooldown drop and an event that never happened were
|
|
// indistinguishable from the operator's side AND from the hub's own records.
|
|
//
|
|
// Measured 2026-08-03: nine `recovery_unit_capture_failed` events arrived, two emails were
|
|
// sent, and the other seven left NO ROW ON ANY CHANNEL. The defect that hid was serious —
|
|
// the cooldown key carries no app identifier, so the first refused app took the slot and
|
|
// every other app's failure that hour was discarded — but the reason it took a day to find
|
|
// the right way round is this line: there was nothing to read.
|
|
//
|
|
// "We chose not to e-mail you" and "nothing happened" must never look identical. This
|
|
// applies to EVERY operator event, not only the one that exposed it. It makes the drop
|
|
// visible; it deliberately does NOT change the cooldown's duration or semantics.
|
|
if err := d.store.LogNotification(customerID, eventType, severity, message,
|
|
"suppressed", "operator cooldown "+window.String()+", key="+cooldownKey, "operator"); err != nil {
|
|
d.logger.Printf("[WARN] Failed to record suppressed operator notification for %s/%s: %v",
|
|
customerID, eventType, err)
|
|
}
|
|
d.logger.Printf("[INFO] Operator email suppressed for %s/%s — cooldown (key=%s)",
|
|
customerID, eventType, cooldownKey)
|
|
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")
|
|
}
|
|
|
|
// recordOnlyEvents are STORED and RECORDED but never e-mailed, on either channel.
|
|
//
|
|
// R-182. The distinction this register exists to make is the whole of that finding: **the record and
|
|
// the notification are different things.** A per-app backup failure must always be written down —
|
|
// every time, unconditionally, regardless of cooldowns, preferences or whether any mail went out —
|
|
// and it must NOT compete for an e-mail slot, because the per-run digest
|
|
// (`backup_run_failures`) is what a person is meant to read.
|
|
//
|
|
// Before this, `recovery_unit_capture_failed` was both at once, and it did neither well: on
|
|
// 2026-08-03 nine of them arrived, two were e-mailed, and the other seven were dropped by the
|
|
// 1-hour cooldown BEFORE anything was written down. So the operator was told about one app, the
|
|
// other apps' failures were discarded, and nothing anywhere recorded that a choice had been made.
|
|
//
|
|
// WHY A REGISTER AND NOT severity "info". Downgrading the severity would have the same routing
|
|
// effect — `severityNotifies` drops info — but it would also relabel a genuine failure as
|
|
// informational in the events table, the operator UI and every historical query, and it would
|
|
// silently drop the X-Priority handling if the type were ever promoted back. This says what it
|
|
// means: not silent, not urgent, RECORDED.
|
|
//
|
|
// IT IS NOT A WAY TO MUTE THINGS. A type belongs here only when something else carries its
|
|
// notification. Adding one with no digest behind it rebuilds the silence R-182 was filed against.
|
|
var recordOnlyEvents = map[string]bool{
|
|
// The per-app Tier-1 capture/refusal failure. Its notification is the run digest, which lists
|
|
// every failed app in one mail; this row is the durable per-failure record behind it.
|
|
"recovery_unit_capture_failed": true,
|
|
}
|
|
|
|
// operatorOnlyEvents are event types that must NEVER reach a customer, whatever their preferences say.
|
|
//
|
|
// R-97c. This register exists because the guarantee it provides was previously ASSERTED IN A COMMENT
|
|
// and not implemented. The claim was that a type with no `customerMessages` entry "structurally
|
|
// cannot" be routed to a customer. It cannot: `FormatCustomerEmail` (templates.go) treats a missing
|
|
// entry as a **fallback to the raw message**, not a block —
|
|
//
|
|
// hunMessage := customerMessages[eventType]
|
|
// if hunMessage == "" { hunMessage = message }
|
|
//
|
|
// — and the only customer gate is `isEventEnabled(prefs.EnabledEvents, ...)`, i.e. CONFIGURATION.
|
|
// So a customer with `whole_guest_backup_failed` in their enabled list and an email set would have
|
|
// received the raw English operator text about a backup they can take no action on.
|
|
//
|
|
// That is the `EffectiveProtected` shape: a doc comment claiming a property the code stopped
|
|
// providing, which is how the samba false alarm survived. The register makes the claim true.
|
|
//
|
|
// NOT implemented as "a missing customerMessages entry blocks delivery" — several existing types rely
|
|
// on the raw-message fallback deliberately (e.g. offbox_enlarge_blocked, whose dynamic Hungarian text
|
|
// is customer-grade and would be DISCARDED by a template). Turning the fallback into a gate would
|
|
// change behaviour well outside this concern.
|
|
var operatorOnlyEvents = map[string]bool{
|
|
// R-97a. A customer can take no action on a failed whole-guest backup, and being told it failed
|
|
// while it is still retrying behind the R-88 breaker is alarming without being actionable.
|
|
"whole_guest_backup_failed": true,
|
|
// The recovery is ALSO listed, even though its customer leg is pairing-gated on a customer-channel
|
|
// "sent" row that can never exist for the line above. Relying on that would make this type's safety
|
|
// a consequence of another type's routing — true today, and silently untrue the moment the failed
|
|
// event becomes customer-visible. Belt, not inference.
|
|
"whole_guest_backup_recovered": true,
|
|
// R-158 / R-167 (D-c). A per-app Tier-1 recovery-unit capture failure. A customer can take no
|
|
// action on it — the causes are a full filesystem, a permission fault or a broken dump, all of
|
|
// which the operator resolves — and the alert carries operator-grade detail (target path, byte
|
|
// figures, the raw error). The customer's half of D-c is the FILL WARNING, which fires BEFORE
|
|
// this and is actionable: free space, delete files, add a drive.
|
|
"recovery_unit_capture_failed": true,
|
|
// R-182. The per-run backup digest. It is the same class as the line above and for the same
|
|
// reason — a customer can act on a full disk (that is the fill warning, which fires first and
|
|
// IS customer-facing) but not on a list of which apps' backups failed and why. It also carries
|
|
// operator-grade detail: per-app leg names, raw refusal reasons and byte figures.
|
|
//
|
|
// Listed here rather than relying on the absence of a `customerMessages` entry, which is NOT a
|
|
// block — `FormatCustomerEmail` falls back to the raw English message. That mistake shipped
|
|
// once (v0.78.0) and the comment above records it.
|
|
"backup_run_failures": true,
|
|
// R-87 (controller v0.231.0). The nightly off-site proof found a backup that is intact and holds
|
|
// none of the app's data. Operator-only for the same reason as `recovery_unit_capture_failed`
|
|
// above: the customer can take no action on it — a hollow recovery unit is a product/ops fault,
|
|
// diagnosed from the capture legs and the manifest, both of which are operator surfaces. The
|
|
// customer's actionable half of this class is the fill warning, which fires first.
|
|
//
|
|
// LISTED IN THE SAME COMMIT THAT MINTS THE TYPE, because a missing customerMessages entry is NOT
|
|
// a block — FormatCustomerEmail falls back to the raw message, which is the v0.78.0 defect this
|
|
// register was built for.
|
|
"offsite_proof_empty": true,
|
|
// R-431 — OPERATOR ONLY. A count that fell is a question for the operator, not a customer:
|
|
// the snapshots still hold the data and telling a customer "your backups were deleted"
|
|
// would be wrong on the usual reading.
|
|
"offsite_snapshots_dropped": true,
|
|
// R-197 (v0.93.0). "The sealed offsite repository key changed" is a custody fact about escrow
|
|
// blobs. A customer can take no action on it — the remedy is the operator's inspection of the
|
|
// off-site tier — and the text is operator-grade English naming host ids and retained-blob
|
|
// counts. Listed here in the SAME commit that mints the type: an operator-tier type that is not
|
|
// registered here reaches customers as raw English, because a missing customerMessages entry is
|
|
// NOT a block (the v0.78.0 defect recorded above).
|
|
"offsite_repo_key_changed": true,
|
|
// R-192 (v0.93.0). These two predate the register and were never added to it — a real gap, not a
|
|
// tidy-up. `offsite_delivery_stuck` is severity warning, has no customerMessages entry, and
|
|
// therefore fell through to FormatCustomerEmail's raw-English fallback: a customer whose box hit
|
|
// the stuck shape was in line for an English e-mail about one-time passwords being "likely
|
|
// burned". Measured on the live hub: notification_log holds operator rows for demo-hp and no
|
|
// customer rows — which is NOT evidence the leg is blocked (it is equally consistent with no
|
|
// configured recipient), so the register makes it structural instead of incidental. Narrowing
|
|
// only: the operator channel is untouched.
|
|
"offsite_delivery_stuck": true,
|
|
"offsite_credential_restaged": true,
|
|
// R-199 (v0.94.0). A host retrieved its own sealed recovery blob. Operator-tier by construction:
|
|
// it names host ids and opaque byte counts, the customer can take no action on it, and its whole
|
|
// purpose is that the operator sees a capability being used. Registered in the same commit that
|
|
// mints the type — an operator-tier type absent from this register reaches customers as raw
|
|
// English (the v0.78.0 defect recorded above).
|
|
"escrow_blob_served": true,
|
|
// R-523 (v0.114.0). The agent restarted a dead in-guest controller, or gave up after repeated
|
|
// restarts. Operator-grade (host ids, vmids, raw docker states) and not actionable by a household
|
|
// — the customer's side of this is the dashboard coming back. Registered in the same commit that
|
|
// mints the types.
|
|
"controller_restarted_by_agent": true,
|
|
"controller_crashloop": true,
|
|
// R-539 (v0.117.0) — the slow sibling: same audience, same reason.
|
|
"controller_slow_crashloop": true,
|
|
// R-518 / R-514 (v0.114.0, controller v0.243.0). A whole-guest tier skipped for absent storage is a
|
|
// provisioning fact the household cannot act on; an OOM-killed app process carries raw container
|
|
// names — the household's side is the dashboard tag. Registered in the same commit.
|
|
"backup_tier_skipped": true,
|
|
"app_oom": true,
|
|
}
|
|
|
|
// IsOperatorOnly reports whether an event type is barred from customer dispatch. Exported so the
|
|
// api package can pin BOTH registers of a new event type in one test — allowlisted-but-not-
|
|
// operator-only is the v0.78.0 defect, and it is only visible when the two are checked together.
|
|
// Read-only: the register itself stays unexported so nothing can widen it at runtime.
|
|
func IsOperatorOnly(eventType string) bool { return operatorOnlyEvents[eventType] }
|
|
|
|
func (d *Dispatcher) processCustomer(customerID, eventType, severity, message, messageCustomer, detailsJSON, source string) {
|
|
// R-97c: operator-tier events stop here, BEFORE prefs are consulted — the point is that no
|
|
// customer configuration can opt in. Logged rather than dropped, so the skip is visible in
|
|
// notification_log instead of looking like a delivery that never happened.
|
|
if operatorOnlyEvents[eventType] {
|
|
d.store.LogNotification(customerID, eventType, severity, message, "skipped", "operator_only", "customer")
|
|
return
|
|
}
|
|
|
|
// 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
|
|
if perAppCustomerCooldownEvents[eventType] {
|
|
cooldownKey += cooldownStackSuffixFor(detailsJSON)
|
|
}
|
|
d.mu.Lock()
|
|
if last, ok := d.custCooldowns[cooldownKey]; ok && time.Since(last) < cooldownDur {
|
|
d.mu.Unlock()
|
|
d.logger.Printf("[INFO] Customer mail skipped for %s/%s — cooldown (key %s)", customerID, eventType, cooldownKey)
|
|
return
|
|
}
|
|
d.custCooldowns[cooldownKey] = time.Now()
|
|
d.mu.Unlock()
|
|
|
|
subject, body := FormatCustomerEmail(d.store.CustomerLanguage(customerID),
|
|
customerID, eventType, severity, message, messageCustomer, 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(d.store.CustomerLanguage(customerID), 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(d.store.CustomerLanguage(customerID), 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
|
|
}
|