Files
felhom.eu/hub/internal/notify/dispatcher.go
T

918 lines
47 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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,
// v0.121.0 (R-636): a storm is one app's event with no digest behind it — two apps storming on one
// night are two alarms. (The controller already sends it at most once per container run.)
"app_oom_storm": true,
// v0.122.0 (R-659): a held app with no whole copy on its box is one app's event with no digest —
// two apps stranded on one night are two alarms, and the second must not be swallowed.
"app_hold_no_whole_copy": true,
// v0.123.0 (decision 28): one app's stop, no digest behind it.
"app_stopped_unhealthy": 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,
// v0.123.0 (decision 28): two apps stopped on one night are two mails.
"app_stopped_unhealthy": 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
}
// perStorageCooldownEvents (R-672, hub v0.124.0): the storage-fill alarms are keyed PER POOL (host +
// storage), with a 6-hour window — "one operator event per pool per 6 hours". The old `customer:type` key
// let one host's pool take the hour and silence another pool filling in it, and a 1-hour window re-mailed
// a pool that sat full for a day. The checker itself emits only on a band ESCALATION, so the window only
// ever collapses a flapping pool.
var perStorageCooldownEvents = map[string]bool{
"storage_fill_warning": true,
"storage_fill_critical": true,
}
const storageFillCooldown = 6 * time.Hour
// cooldownStorageSuffix returns ":"+host_id+"/"+storage for a per-storage type, else "".
func cooldownStorageSuffix(eventType, detailsJSON string) string {
if !perStorageCooldownEvents[eventType] || detailsJSON == "" {
return ""
}
var d struct {
HostID string `json:"host_id"`
Storage string `json:"storage"`
}
if err := json.Unmarshal([]byte(detailsJSON), &d); err != nil || d.Storage == "" {
return ""
}
return ":" + d.HostID + "/" + d.Storage
}
// 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
}
if perStorageCooldownEvents[eventType] {
return storageFillCooldown
}
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) + cooldownStorageSuffix(eventType, detailsJSON)
window := operatorCooldownFor(eventType)
// R-723 (v0.126.0): a whole-guest tier skipped for absent storage in a box's FIRST HOUR is provisioning
// still in flight, not a fault (measured 2026-09-29: the first local backup ran 7 min after enrolment,
// the off-site descriptor arrived with the next 15-min host report, the off-site tier then backed up
// fine — and the operator had been mailed „skipped"). Recorded, not mailed.
if eventType == "backup_tier_skipped" && d.hostEnrolledWithin(customerID, firstHourOfABox) {
if err := d.store.LogNotification(customerID, eventType, severity, message,
"suppressed", "first hour of a new box (R-723)", "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 — the box enrolled less than %s ago", customerID, eventType, firstHourOfABox)
return
}
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,
// v0.127.0 (decisions 68–69, R-820/R-822). The off-site key registrar, its daily check and the
// clean-up window: custody facts about key lines and fingerprints — the household can take no
// action on any of them. Listed in the SAME commit that mints them.
"offsite_key_installed": true,
"offsite_key_unlocked": true,
"offsite_key_audit_failed": true,
"offsite_repo_moved_aside": true,
"offsite_window_drop": true,
"offsite_window_failed": true,
"offsite_prune_guard_refused": true,
// R-833 (v0.129.0): the operator raised ONE window's removal cap — an operator act, logged.
"offsite_window_large_grant": true,
// OS updates (hub v0.130.0, `11` §8 step 2): run failures, rings, switches and approvals are operator facts.
// os_update_applied is deliberately NOT here — it is the household's one line (info: recorded, never mailed).
"os_update_failed": true,
"os_update_health_failed": true,
"os_release_approved": true,
"os_release_approved_now": true,
"os_update_settings_changed": true,
// R-841 (hub v0.131.0): the tunnel alarm — a box fact the household can do nothing about from inside.
"tunnel_down": true,
"tunnel_recovered": true,
// `11` §8.3 (hub v0.131.0): the four OS-update alarms — fleet facts only the operator can act on.
"os_update_stale": true,
"os_reboot_needed": true,
"os_ring0_stalled": true,
"os_not_covered": true,
// R-851 (hub v0.132.0): the crash guard. host_restarted_after_crash is deliberately NOT here — the household's line.
"host_crash_restart": true,
"host_crash_guard_tripped": true,
"host_crash_guard_rearmed": true,
"host_kernel_oops": 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,
// R-636 (v0.121.0, controller v0.265.0). The LOUDER sibling of app_oom: the same container run
// OOM-killed 20+ times in 30 minutes. Same audience, same reason — raw container names and memory
// figures; the household's side is the dashboard. Registered in the same commit that mints it.
"app_oom_storm": true,
// R-659 (v0.122.0, controller v0.268.0; operator ruling 2026-09-24, `09` §3 decision 25). A held app
// whose box holds no copy that brings it back WHOLE. The household is told by its own
// `app_update_held` mail, in its language, that support is informed; THIS is that information —
// operator-grade (the copies seen, per tier, with dates), and the act it calls for is support's.
// Registered in the same commit that mints it.
"app_hold_no_whole_copy": 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
}
// firstHourOfABox is how long after enrolment a whole-guest tier skip is provisioning, not a fault (R-723).
const firstHourOfABox = time.Hour
// hostEnrolledWithin reports whether the customer's current host was enrolled less than d ago.
func (d *Dispatcher) hostEnrolledWithin(customerID string, within time.Duration) bool {
if d.store == nil {
return false
}
h, err := d.store.GetHostByCustomer(customerID)
if err != nil || h == nil || h.CreatedAt.IsZero() {
return false
}
return time.Since(h.CreatedAt) < within
}