Files
felhom-controller/controller/internal/notify/notifier.go
T
admin bef39598d0
gates / gates (push) Successful in 24s
v0.256.0: the box sends its own sentence in the household's language (R-558 Part B)
MinAgent: 0.131.0 (unchanged). Needs hub v0.118.0+, which shipped first and
tolerates a box that sends none of this - every box in the fleet is that box
until this release reaches it.

The hub writes a household's e-mails in their language now, but about a third
of those mails carry a sentence the BOX composed, naming a drive, an app or a
number. The hub cannot translate one. So the box sends it twice.

- message_customer on POST /api/v1/event, omitempty. A HUNGARIAN household
  sends nothing extra at all, so its payload stays byte-for-byte what every box
  sends today and the hub's fallback path keeps being the one production
  exercises rather than a branch nobody takes.
- 19 producers render both sentences from ONE bundle key. `message` stays
  Hungarian always: it is what the operator is mailed and what the hub logs.
- customer.language bootstraps a new box - stored choice, then config, then
  Hungarian. The config value is NEVER written into settings.json: that would
  record a choice the household never made.

The Hungarian did not move, measured twice: the wire golden from the slice-2
base commit, and the Go parity gate over all 19 new keys.

Three guards had to learn the change and one caught me: the test seam now
carries the new field; the R-329 severity register reported two dynamic sites
as no longer existing the moment they moved off PushEvent (the walk now checks
36 severity literals, up from 20); and TestConfigLanguageIsWiredInMain reads
main.go, because cmd/ is gitignored and ripgrep does not.

A mistake, named: the first pass dropped displayName from three producers,
which would have mailed customers "Alkalmazás telepítve: %!s(MISSING)". Caught
reading the diff; now pinned by a test that refuses %!/MISSING/%s/%d in either
language.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
2026-09-18 16:54:20 +02:00

1131 lines
47 KiB
Go

package notify
import (
"bytes"
"encoding/json"
"fmt"
"gitea.dooplex.hu/admin/felhom-controller/internal/util"
"io"
"log"
"net/http"
"strings"
"sync"
"time"
"gitea.dooplex.hu/admin/felhom-controller/internal/settings"
"gitea.dooplex.hu/admin/felhom-controller/internal/i18n"
)
// Notifier sends structured events to the hub via /api/v1/event.
// Non-blocking: fires requests in goroutines, logs errors but doesn't retry aggressively.
// Cooldown logic is handled by the Hub — the controller sends all events unconditionally.
// EventHistoryEntry records a sent event for the debug page.
type EventHistoryEntry struct {
Timestamp time.Time `json:"timestamp"`
EventType string `json:"event_type"`
Severity string `json:"severity"`
Message string `json:"message"`
HubStatus int `json:"hub_status"`
HubError string `json:"hub_error,omitempty"`
}
type Notifier struct {
hubURL string
apiKey string
customerID string
httpClient *http.Client
logger *log.Logger
enabled bool
debug bool
settings *settings.Settings
mu sync.Mutex
prevHealthStatus string // tracks previous health check status for change detection
// oomSeen (R-514) remembers container runs already reported as OOM-killed.
oomSeen map[string]bool
// appDown tracks which deployed apps are currently in the DOWN state so app_start_failed fires
// ONCE per running→down transition, not every health cycle (fix-3 anti-spam). In-memory: a
// controller restart re-notifies once (acceptable — better than missing). The hub owns the real
// cooldown; the controller must not add its own timer.
appDown map[string]bool
// pushFn is a test seam for the transition-emitting notifiers (fix-3). nil → the real async
// PushEvent; tests inject a synchronous recorder.
// v0.256.0 (R-558): messageCustomer is IN the seam, not beside it. A seam that cannot see a new
// field is a seam that cannot test it — the same lesson this file already records one comment
// down, where PushEvent had to become the seam because emit alone made a whole class of
// producer invisible.
pushFn func(eventType, severity, message, messageCustomer string, details interface{})
// Event history ring buffer (debug page)
historyMu sync.RWMutex
history [50]EventHistoryEntry
histPos int
histFull bool
}
// New creates a new Notifier. Returns a no-op notifier if hub is not enabled.
func New(hubURL, apiKey, customerID string, sett *settings.Settings, logger *log.Logger, debug bool) *Notifier {
enabled := hubURL != "" && apiKey != ""
if enabled {
logger.Printf("[INFO] Notifier enabled (hub: %s)", hubURL)
} else {
logger.Printf("[INFO] Notifier disabled (hub not configured)")
}
return &Notifier{
hubURL: hubURL,
apiKey: apiKey,
customerID: customerID,
httpClient: &http.Client{Timeout: 10 * time.Second},
logger: logger,
enabled: enabled,
debug: debug,
settings: sett,
}
}
// IsEnabled returns whether the notifier has a configured hub connection.
func (n *Notifier) IsEnabled() bool {
return n.enabled
}
// ── Detail structs ───────────────────────────────────────────────────
// BackupDetails holds structured data for backup events.
type BackupDetails struct {
DriveCount int `json:"drive_count,omitempty"`
SnapshotID string `json:"snapshot_id,omitempty"`
DurationSec int `json:"duration_sec,omitempty"`
DataAdded string `json:"data_added,omitempty"`
Error string `json:"error,omitempty"`
}
// WholeGuestBackupDetails carries the TIER for a whole-guest (vzdump) backup event (R-97a).
//
// The `tier` field is load-bearing beyond display: the hub's operator cooldown is keyed
// `customerID:eventType` plus this tier when present, so `local` failing does not get swallowed by
// `felhom-pbs` having failed within the same hour. Rename it and the two tiers silently share one
// cooldown again.
type WholeGuestBackupDetails struct {
Tier string `json:"tier"`
Error string `json:"error,omitempty"`
}
// DBDumpDetails holds structured data for DB dump events.
type DBDumpDetails struct {
DatabaseCount int `json:"database_count,omitempty"`
TotalSize string `json:"total_size,omitempty"`
DurationSec int `json:"duration_sec,omitempty"`
Error string `json:"error,omitempty"`
}
// DiskDetails holds structured data for disk warning/critical events.
type DiskDetails struct {
Mount string `json:"mount,omitempty"`
UsagePercent float64 `json:"usage_percent,omitempty"`
Label string `json:"label,omitempty"`
}
// HealthDetails holds structured data for health events.
type HealthDetails struct {
PreviousStatus string `json:"previous_status,omitempty"`
CurrentStatus string `json:"current_status,omitempty"`
Issues []string `json:"issues,omitempty"`
Warnings []string `json:"warnings,omitempty"`
}
// StorageDetails holds structured data for storage events.
type StorageDetails struct {
DrivePath string `json:"drive_path,omitempty"`
Label string `json:"label,omitempty"`
StoppedApps []string `json:"stopped_apps,omitempty"`
}
// UpdateDetails holds structured data for controller update events.
type UpdateDetails struct {
FromVersion string `json:"from_version,omitempty"`
ToVersion string `json:"to_version,omitempty"`
Error string `json:"error,omitempty"`
}
// AppDetails holds structured data for app lifecycle events.
type AppDetails struct {
StackName string `json:"stack_name,omitempty"`
DisplayName string `json:"display_name,omitempty"`
}
// CrossDriveDetails holds structured data for cross-drive backup events.
type CrossDriveDetails struct {
StackName string `json:"stack_name,omitempty"`
Method string `json:"method,omitempty"`
DestPath string `json:"dest_path,omitempty"`
Duration string `json:"duration,omitempty"`
Error string `json:"error,omitempty"`
}
// ── Core event push ──────────────────────────────────────────────────
// eventRequest is the JSON payload sent to /api/v1/event.
type eventRequest struct {
CustomerID string `json:"customer_id"`
EventType string `json:"event_type"`
Severity string `json:"severity"`
Message string `json:"message"`
// MessageCustomer (v0.256.0, R-558) is the SAME sentence in the household's language.
//
// `omitempty` is load-bearing: a Hungarian household sends no second copy at all, so its event
// payload is byte-for-byte what it has always been, and the hub's fallback path — the one every
// box in the fleet uses today — is the one that keeps being exercised in production rather than
// becoming a branch nobody takes.
//
// Message stays Hungarian ALWAYS. It is what the operator is mailed, what the hub logs and what
// notification_log records; nothing an operator reads may move because a household switched.
MessageCustomer string `json:"message_customer,omitempty"`
Details json.RawMessage `json:"details,omitempty"`
}
// PushEvent sends a structured event to the hub's /api/v1/event endpoint.
// Non-blocking (goroutine). Retries twice with 3s backoff.
// details may be nil (omitted from JSON) or a struct that marshals to JSON.
// pushEventMsg renders ONE bundle key twice — Hungarian for the wire, the household's language
// beside it — and pushes both (R-558).
//
// Producers call this instead of composing a Hungarian sentence, because the hub CANNOT translate a
// sentence the box composed: it arrives as finished text naming a drive, an app or a number. The
// only place both languages can be produced is here, where the key and its arguments still exist.
//
// The Hungarian is rendered from the SAME key, so it cannot drift from the English: there is one
// sentence with two spellings, not two sentences. That the Hungarian is byte-identical to the
// literal it replaced is pinned by TestEventMessageWireTextIsFrozen and by the Go parity gate.
func (n *Notifier) pushEventMsg(eventType, severity, key string, details interface{}, args ...interface{}) {
b, err := i18n.Shared()
if err != nil {
// The bundle is embedded, so this is a broken build rather than a runtime condition. Push
// the key itself rather than an empty sentence: an event with a visible key reaches the
// operator and gets fixed; an event with an empty message is a silent hole.
n.logger.Printf("[ERROR] pushEventMsg: bundle unavailable for %s: %v", key, err)
n.pushEventBoth(eventType, severity, key, "", details)
return
}
hungarian := b.Msgf(i18n.Default, key, args...)
// A Hungarian household sends nothing extra — see MessageCustomer's comment.
household := ""
if lang := n.boxLang(); lang != i18n.Default {
household = b.Msgf(lang, key, args...)
}
n.pushEventBoth(eventType, severity, hungarian, household, details)
}
// pushEventMsgSuffix is pushEventMsg with a tail that CANNOT be translated — a docker error, a
// validator's sentence, something that arrived as finished text. It is appended to both renderings,
// so an English household gets an English sentence with a foreign tail rather than no sentence.
func (n *Notifier) pushEventMsgSuffix(eventType, severity, key, suffix string, details interface{}, args ...interface{}) {
b, err := i18n.Shared()
if err != nil {
n.logger.Printf("[ERROR] pushEventMsgSuffix: bundle unavailable for %s: %v", key, err)
n.pushEventBoth(eventType, severity, key+suffix, "", details)
return
}
household := ""
if lang := n.boxLang(); lang != i18n.Default {
household = b.Msgf(lang, key, args...) + suffix
}
n.pushEventBoth(eventType, severity, b.Msgf(i18n.Default, key, args...)+suffix, household, details)
}
// boxLang is the household's language, nil-safe. A notifier built without settings (several tests,
// and the setup-mode path) is Hungarian rather than a panic.
func (n *Notifier) boxLang() string {
if n.settings == nil {
return i18n.Default
}
return n.settings.GetLanguage()
}
// PushEvent sends a structured event to the hub's /api/v1/event endpoint, Hungarian only.
//
// Kept for producers whose sentence is composed somewhere else and reaches them already finished —
// and for the OPERATOR-tier types, which are Hungarian (or English) by design and have no household
// half. A customer-facing producer should use pushEventMsg.
func (n *Notifier) PushEvent(eventType, severity, message string, details interface{}) {
n.pushEventBoth(eventType, severity, message, "", details)
}
func (n *Notifier) pushEventBoth(eventType, severity, message, messageCustomer string, details interface{}) {
// The test seam sits HERE, not only in emit (v0.252.0, R-557). Most producers call PushEvent
// directly, so a test that set pushFn saw nothing from them and the difference was invisible: a
// producer that pushes nothing and a producer whose message is empty both leave a recorder empty.
// TestEventMessageWireTextIsFrozen pins the exact bytes the hub mails to a household, and it can
// only do that if every producer is reachable through one seam. In production pushFn is nil and
// this line does nothing.
if n.pushFn != nil {
n.pushFn(eventType, severity, message, messageCustomer, details)
return
}
if !n.enabled {
return
}
var detailsJSON json.RawMessage
if details != nil {
b, err := json.Marshal(details)
if err != nil {
n.logger.Printf("[WARN] PushEvent: failed to marshal details for %s: %v", eventType, err)
} else {
detailsJSON = b
}
}
payload := eventRequest{
CustomerID: n.customerID,
EventType: eventType,
Severity: severity,
Message: message,
MessageCustomer: messageCustomer,
Details: detailsJSON,
}
jsonData, err := json.Marshal(payload)
if err != nil {
n.logger.Printf("[ERROR] PushEvent: marshal failed for %s: %v", eventType, err)
return
}
go func() {
url := n.hubURL + "/api/v1/event"
if n.debug {
n.logger.Printf("[DEBUG] PushEvent: type=%s severity=%s url=%s", eventType, severity, url)
}
var lastErr error
for attempt := 0; attempt < 3; attempt++ {
if attempt > 0 {
time.Sleep(3 * time.Second)
}
req, err := http.NewRequest("POST", url, bytes.NewReader(jsonData))
if err != nil {
lastErr = err
continue
}
req.Header.Set("Authorization", "Bearer "+n.apiKey)
req.Header.Set("Content-Type", "application/json")
resp, err := n.httpClient.Do(req)
if err != nil {
lastErr = err
continue
}
io.Copy(io.Discard, resp.Body)
resp.Body.Close()
if resp.StatusCode >= 200 && resp.StatusCode < 300 {
if n.debug {
n.logger.Printf("[DEBUG] PushEvent: %s pushed OK (HTTP %d)", eventType, resp.StatusCode)
}
n.logger.Printf("[INFO] Event pushed: %s (%s) — %s", eventType, severity, message)
n.recordHistory(eventType, severity, message, resp.StatusCode, "")
return
}
lastErr = fmt.Errorf("HTTP %d", resp.StatusCode)
}
n.logger.Printf("[WARN] Event push failed after 3 attempts (%s/%s): %v", eventType, severity, lastErr)
errMsg := ""
if lastErr != nil {
errMsg = lastErr.Error()
}
n.recordHistory(eventType, severity, message, 0, errMsg)
}()
}
// ── Convenience methods ──────────────────────────────────────────────
// NotifyHealthChange checks if health status changed and sends appropriate events.
// Detects both degradation (ok→warn, ok→fail, warn→fail) and recovery (fail→ok, warn→ok, fail→warn).
func (n *Notifier) NotifyHealthChange(status string, issues, warnings []string) {
if !n.enabled {
return
}
n.mu.Lock()
prev := n.prevHealthStatus
n.prevHealthStatus = status
n.mu.Unlock()
if prev == "" {
return // First run, just record status
}
if status == prev {
return
}
details := HealthDetails{
PreviousStatus: prev,
CurrentStatus: status,
Issues: issues,
Warnings: warnings,
}
prevRank := statusRank(prev)
newRank := statusRank(status)
if newRank > prevRank {
// Degradation
if status == "fail" {
n.pushEventMsg("health_critical", "error", "event.health_critical", details, prev)
} else if status == "warn" {
n.pushEventMsg("health_degraded", "warning", "event.health_degraded", details, prev)
}
} else {
// Recovery
n.pushEventMsg("health_recovered", "info", "event.health_recovered", details, status, prev)
}
}
// NotifyBackupFailed sends a backup failure event.
func (n *Notifier) NotifyBackupFailed(message, errMsg string) {
n.PushEvent("backup_failed", "error", message, BackupDetails{Error: errMsg})
}
// RecoveryUnitFailureDetails is the machine-readable tail of a Tier-1 capture failure. App NAMES and
// byte figures only — never an env value (§9.5).
type RecoveryUnitFailureDetails struct {
App string `json:"app"`
Error string `json:"error"`
TargetPath string `json:"target_path,omitempty"`
UsedGB float64 `json:"used_gb,omitempty"`
AvailGB float64 `json:"avail_gb,omitempty"`
TotalGB float64 `json:"total_gb,omitempty"`
UsedPercent float64 `json:"used_percent,omitempty"`
// SpaceKnown distinguishes "we read the filesystem and it says these numbers" from "we could not
// read it". Without it, an unreadable target is indistinguishable from an empty one — the
// presence-is-not-success trap, in the other direction.
SpaceKnown bool `json:"space_known"`
}
// NotifyRecoveryUnitCaptureFailed sends the OPERATOR-TIER alert for a per-app Tier-1 recovery-unit
// capture failure (R-158, D-c's operator half).
//
// DELIBERATELY NOT `backup_failed`. That type carries a `customerMessages` entry AND sits in
// `settings.DefaultEnabledEvents`, so reusing it would email the customer, in Hungarian, that their
// backup failed — an event they can take no action on. It is exactly the mistake R-97a avoided by
// minting `whole_guest_backup_failed`, and the reasoning is written into the hub's handler.go.
// R-158's original proposal named `backup_failed`; decision D-c routes this to the operator, and
// where the two disagree D-c wins.
//
// Operator-only is enforced by the hub's `notify.operatorOnlyEvents` register, NOT by the absence of
// a customerMessages entry — v0.78.0 claimed the latter and was wrong.
func (n *Notifier) NotifyRecoveryUnitCaptureFailed(message string, d RecoveryUnitFailureDetails) {
n.PushEvent("recovery_unit_capture_failed", "error", message, d)
}
// RunFailureDetail is one app's failed leg inside a backup run digest.
type RunFailureDetail struct {
App string `json:"app"`
Leg string `json:"leg"`
Reason string `json:"reason"`
}
// BackupRunFailuresDetails is the per-RUN digest payload (R-182). App NAMES, leg names, reasons and
// byte figures only — never an env value (§9.4).
type BackupRunFailuresDetails struct {
// RunID makes the hub's 1-hour operator cooldown unable to collapse two real runs into one
// e-mail. EMPTY on the periodic refresh sweep, deliberately: that path can fire on every status
// poll, so it must fall under the ordinary cooldown instead.
RunID string `json:"run_id,omitempty"`
RunKind string `json:"run_kind"`
Failed int `json:"failed"`
Attempted int `json:"attempted"`
TargetPath string `json:"target_path,omitempty"`
UsedGB float64 `json:"used_gb,omitempty"`
AvailGB float64 `json:"avail_gb,omitempty"`
TotalGB float64 `json:"total_gb,omitempty"`
UsedPercent float64 `json:"used_percent,omitempty"`
SpaceKnown bool `json:"space_known"`
Apps []RunFailureDetail `json:"apps"`
}
// NotifyBackupRunFailures sends the ONE operator digest for a backup run in which something failed
// (R-182). It is the NOTIFICATION; the per-app `recovery_unit_capture_failed` events are the RECORD,
// and the hub routes those record-only so they never compete for an e-mail slot.
//
// OPERATOR-TIER, and for the same reason as its per-app sibling: 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. `notify.operatorOnlyEvents` in the hub enforces that — NOT the
// absence of a customerMessages entry, which is a fallback rather than a block (v0.78.0).
func (n *Notifier) NotifyBackupRunFailures(message string, d BackupRunFailuresDetails) {
n.PushEvent("backup_run_failures", "error", message, d)
}
// NotifyOffboxEnlargeBlocked sends a WARNING (not a failure) when an app's enlarged offsite push was
// refused by the pre-push quota gate — its config+DB were still saved. Customer-facing (Hungarian
// body). NOTE: the event type "offbox_enlarge_blocked" must be added to the hub's allowedEventTypes +
// customerMessages for delivery (a hub-side task, flagged — until then the hub 400s/drops it and the
// in-dashboard LastWarning + /backups/remote note carry the message).
func (n *Notifier) NotifyOffboxEnlargeBlocked(message string) {
n.PushEvent("offbox_enlarge_blocked", "warning", message, nil)
}
// (NotifyBackupCompleted removed 2026-06-16 — the backup_completed event had no callers
// since slice 8C moved whole-guest backup to the agent. The hub's backup-deadline check
// now reads the agent host-report's PBS snapshots instead of this event. DB-dump events
// below are still emitted and consumed.)
// NotifyDBDumpFailed sends a DB dump failure event.
func (n *Notifier) NotifyDBDumpFailed(message, errMsg string) {
n.PushEvent("db_dump_failed", "error", message, DBDumpDetails{Error: errMsg})
}
// NotifyDBDumpCompleted sends a DB dump success event.
func (n *Notifier) NotifyDBDumpCompleted(details DBDumpDetails) {
n.pushEventMsg("db_dump_completed", "info", "event.db_dump_completed", details)
}
// NotifyIntegrityFailed sends a backup integrity check failure event.
func (n *Notifier) NotifyIntegrityFailed(message, errMsg string) {
n.PushEvent("backup_integrity_failed", "error", message, &BackupDetails{Error: errMsg})
}
// NotifyOffsiteProofEmpty (R-87) reports that the nightly off-site proof found a backup that is
// READABLE and holds none of the app's data.
//
// SEVERITY `error`, from the hub's exact vocabulary {info, warning, error, critical}. Anything else
// is silently coerced to `info` and mailed to nobody — that shipped twice (R-328 on
// `disk_health_degraded`, R-329 on `app_start_failed`, 91 events stored and zero delivered), and
// `r329_severity_contract_test.go` walks this file to keep it from shipping a third time.
//
// A SEPARATE TYPE FROM `backup_integrity_failed`, and that is the point rather than an oversight: the
// integrity check says THE STORE IS DAMAGED; this says the store is sound and the CONTENT is absent.
// The customer's action differs and so must the sentence.
//
// The message is passed through, not templated hub-side, so it can name the app and what is missing.
func (n *Notifier) NotifyOffsiteProofEmpty(message, detail string) {
n.PushEvent("offsite_proof_empty", "error", message, &BackupDetails{Error: detail})
}
// NotifyIntegrityOK sends a backup integrity check success event.
func (n *Notifier) NotifyIntegrityOK(message string) {
n.PushEvent("backup_integrity_ok", "info", message, nil)
}
// NotifyControllerUpdated sends a controller update event.
func (n *Notifier) NotifyControllerUpdated(fromVer, toVer string, success bool) {
severity := "info"
key := "event.controller_updated"
details := UpdateDetails{FromVersion: fromVer, ToVersion: toVer}
if !success {
severity = "error"
key = "event.controller_update_failed"
}
n.pushEventMsg("controller_updated", severity, key, details, fromVer, toVer)
}
// NotifyControllerStarted sends a controller startup event.
// details may include self-test summary (e.g., {"selftest_pass": 8, "selftest_warn": 1, "selftest_fail": 0}).
func (n *Notifier) NotifyControllerStarted(version string, details map[string]interface{}) {
n.pushEventMsg("controller_started", "info", "event.controller_started", details, version)
}
// NotifyStorageDisconnected sends a drive disconnection event.
func (n *Notifier) NotifyStorageDisconnected(label string, stoppedApps []string) {
n.pushEventMsg("storage_disconnected", "error", "event.storage_disconnected", StorageDetails{
Label: label,
StoppedApps: stoppedApps,
}, label)
}
// NotifyBackupTargetAbsent (E-2) reports that the drive holding the WHOLE-GUEST backup is gone.
//
// Distinct from NotifyStorageDisconnected on purpose. That one means "a drive went away and some apps
// may have stopped"; this means "the thing that makes your backup survive a disk failure is gone" —
// a different customer action and a different operator urgency. Before E-2 this had NO prompt signal
// at all: the tier stays DUE (targetStoragePresent checks name presence, never reachability), so the
// only evidence was its own failure at the next due cycle, up to ~24 h away on the daily local tier.
func (n *Notifier) NotifyBackupTargetAbsent(label, target string) {
n.pushEventMsg("backup_target_absent", "error", "event.backup_target_absent",
StorageDetails{Label: label}, label, target)
}
// NotifyBackupTargetRestored is the paired recovery. info severity — the existing recovery pattern;
// severityNotifies is deliberately NOT widened.
func (n *Notifier) NotifyBackupTargetRestored(label, target string) {
n.pushEventMsg("backup_target_restored", "info", "event.backup_target_restored",
StorageDetails{Label: label}, label, target)
}
// NotifyStorageReconnected sends a drive reconnection event.
func (n *Notifier) NotifyStorageReconnected(label string) {
n.pushEventMsg("storage_reconnected", "info", "event.storage_reconnected",
StorageDetails{Label: label}, label)
}
// AgentChannelDetails carries the classified reason for a controller→agent channel-down event.
type AgentChannelDetails struct {
Reason string `json:"reason"`
}
// NotifyAgentChannelDown sends an OPERATOR-facing controller→agent channel-down event (the message is
// English — the customer can't act on "the agent re-keyed", so this never reaches them: the event type
// is not a customer notification toggle, same as the host_* operator events). eventType/severity/
// message come from the channelhealth classifier (spike Q1 map).
func (n *Notifier) NotifyAgentChannelDown(reason, eventType, severity, message string) {
n.PushEvent(eventType, severity, message, AgentChannelDetails{Reason: reason})
}
// NotifyAgentChannelRecovered sends the recovery event (info → logged, no email, mirrors
// storage_reconnected).
func (n *Notifier) NotifyAgentChannelRecovered() {
n.PushEvent("agent_channel_recovered", "info",
"Controller→agent channel recovered — local-API reachable again.", nil)
}
// NotifyEndpointDrift reports a local_api endpoint divergence (R-77). Operator-only English, its
// OWN event type — deliberately not folded into agent_channel_*, because during the 2026-07-25
// outage the generic channel alert was the only signal and it hid a specific, fixable config fault.
// severity=error: unlike a transient unreachable, drift never self-heals.
//
// NOTE: the hub validates event_type against allowedEventTypes and 400s an unknown one, so this
// type MUST exist there too (hub handler.go) or the alert is silently inert.
func (n *Notifier) NotifyEndpointDrift(message string, fingerprintAgrees bool) {
n.PushEvent("local_api_endpoint_drift", "error", message,
EndpointDriftDetails{FingerprintAgrees: fingerprintAgrees})
}
// EndpointDriftDetails carries NO addresses and NO secrets — the endpoints are in the message, and
// the pin is a boolean by design.
type EndpointDriftDetails struct {
FingerprintAgrees bool `json:"fingerprint_agrees"`
}
// NotifyAppDeployed sends an app deployment event.
//
// R-536: it is called when the deploy COMPLETES (the compose up succeeded and the app's durable
// state was written), never when the request was merely accepted. It used to fire beside the 202,
// so an install that was interrupted five seconds later still stood on the hub's timeline as
// „Alkalmazás telepítve" forever — measured 2026-09-16 with mealie, which ended `not_deployed`.
func (n *Notifier) NotifyAppDeployed(stackName, displayName string) {
n.pushEventMsg("app_deployed", "info", "event.app_deployed",
AppDetails{StackName: stackName, DisplayName: displayName}, displayName)
}
// NotifyAppDeployStarted records the ACCEPTANCE — the fact `app_deployed` used to assert. It keeps
// the timeline's "the customer asked for this app at 12:31" without claiming the install finished.
//
// NOTE: the hub validates event_type against allowedEventTypes and 400s an unknown one, so this type
// MUST exist there too (hub handler.go) or the event is silently inert.
func (n *Notifier) NotifyAppDeployStarted(stackName, displayName string) {
n.pushEventMsg("app_deploy_started", "info", "event.app_deploy_started",
AppDetails{StackName: stackName, DisplayName: displayName}, displayName)
}
// NotifyAppDeployFailed closes the pair. severity=warning, not info: an install the customer started
// and that did not finish is a thing someone should see, and the silent version of this is exactly
// what left a completed-install record for an app that was never installed.
func (n *Notifier) NotifyAppDeployFailed(stackName, displayName, reason string) {
// The reason is appended to BOTH renderings rather than folded into the key: it arrives as
// finished text (a docker error, a validator's sentence) that this function cannot translate.
suffix := ""
if reason != "" {
suffix = " — " + reason
}
n.pushEventMsgSuffix("app_deploy_failed", "warning", "event.app_deploy_failed", suffix,
AppDetails{StackName: stackName, DisplayName: displayName}, displayName)
}
// AppRunState is one deployed app's running state for the fix-3 start-failure notifier: Down=true
// when the app is deployed but its containers are not running.
type AppRunState struct {
Name string
DisplayName string
Down bool
// IntentUnknown is set when this app is STOPPED and no customer intent was ever recorded, so the
// alarm was suppressed by the §4 fallback rather than by a decision anyone made (R-386).
//
// It rides here rather than being a third return value or a logger parameter so that the log line
// and the suppression come from the SAME computation — a separately-derived log is a second
// source of truth, and the two drift. `classifyRunStates` stays pure and its signature does not
// move, which is what lets its existing tests keep testing what they were written to test.
IntentUnknown bool
}
// NotifyAppStartFailures fires an `app_start_failed` hub event ONCE per running→down transition
// (fix-3). It is called each health cycle with the CURRENT deployed-app run states; the per-app
// transition tracking (n.appDown) makes down→down cycles silent, so a persistently-dead app does not
// spam. down→running clears the tracker (no event — the dashboard banner self-clears; a recovery
// event is deliberately omitted to keep the operator inbox quiet). The hub applies its own cooldown.
func (n *Notifier) NotifyAppStartFailures(apps []AppRunState) {
n.mu.Lock()
if n.appDown == nil {
n.appDown = map[string]bool{}
}
var newlyDown []AppRunState
seen := map[string]bool{}
for _, a := range apps {
seen[a.Name] = true
was := n.appDown[a.Name]
if a.Down && !was {
newlyDown = append(newlyDown, a) // running→down (or first-seen-down after the boot grace)
}
n.appDown[a.Name] = a.Down
}
// Forget apps no longer reported (removed/undeployed) so a later redeploy re-notifies cleanly.
for name := range n.appDown {
if !seen[name] {
delete(n.appDown, name)
}
}
n.mu.Unlock()
for _, a := range newlyDown {
name := a.DisplayName
if name == "" {
name = a.Name
}
// R-329: "warning", NOT "warn". The hub's vocabulary is exactly
// {info, warning, error, critical} and it COERCES anything else to "info" at ingest, silently
// — after which severityNotifies drops it and NEITHER leg runs. See DiskAlertKind.Severity's
// doc comment, which records the same mistake shipping once before (v0.215.0). This one was
// worse: it was invisible for months because R-384's ordering defect meant the event could
// not fire at all, so a broken severity had nothing to break.
// Pinned by TestR329_EveryEmittedSeverityIsInTheHubVocabulary (AST walk over this package).
n.emit("app_start_failed", "warning",
fmt.Sprintf("Telepített alkalmazás nem fut: %s", name),
AppDetails{StackName: a.Name, DisplayName: a.DisplayName})
}
}
// NotifyAppOOM — R-514 (v0.243.0). A process inside a running app container was killed by the memory
// limit (Docker State.OOMKilled). Fires ONCE per container run (keyed by StartedAt), so a container
// that stays OOM-marked does not repeat. "warning" — the hub vocabulary (R-329). Operator-only
// hub-side (hub >= v0.114.0 registers it): the household sees the dashboard tag.
func (n *Notifier) NotifyAppOOM(stack, container, startedAt string) {
key := container + "|" + startedAt
n.mu.Lock()
if n.oomSeen == nil {
n.oomSeen = map[string]bool{}
}
if n.oomSeen[key] {
n.mu.Unlock()
return
}
n.oomSeen[key] = true
n.mu.Unlock()
n.emit("app_oom", "warning",
fmt.Sprintf("Alkalmazás memóriája elfogyott: %s (%s) — egy folyamatát a memóriakorlát leállította", stack, container),
AppDetails{StackName: stack, DisplayName: container})
}
// DiskHealthDetails is the event-detail payload for disk_health_degraded.
type DiskHealthDetails struct {
Disk string `json:"disk"`
Attributes []string `json:"attributes,omitempty"`
Critical bool `json:"critical"`
}
// DiskAlertKind selects the customer-facing message shape. A customer's ACTION differs by kind —
// "back up and call us for a replacement" is not "check the ventilation" — so the kind travels with
// the alert rather than being flattened into a single sentence.
type DiskAlertKind int
const (
DiskAlertWarn DiskAlertKind = iota // Figyelmeztetés — worth keeping an eye on
DiskAlertFailSelfReported // Hiba — the drive's own SMART verdict says FAILING
DiskAlertFailSectors // Hiba — reached from unreadable-sector counters
DiskAlertFailTemperature // Hiba — reached from heat
DiskAlertFailWorsened // Hiba — already reported, and still getting worse
)
// DiskAlert is the payload for one disk-health alert. It carries enough for the notifier to pick a
// message shape and fill in the counts; message CONSTRUCTION stays here because the notifier owns
// customer copy, and moving it to the caller would scatter Hungarian across packages.
type DiskAlert struct {
Label string // customer-facing disk label (device model where known)
Attributes []string // Hungarian attribute names behind the verdict (nil for a self-reported FAILING)
Kind DiskAlertKind
Sectors int // max(pending, offline_uncorrectable) — quoted in the sector/worsened shapes
TemperatureC int // °C — quoted in the temperature shape
}
// Severity is the hub-accepted severity string for this alert.
//
// THE VOCABULARY IS EXACT AND IT IS THE HUB'S, NOT OURS. The hub accepts only
// {"info","warning","error","critical"} and silently COERCES anything else to "info"
// (felhom.eu/hub/internal/api/handler.go, the severity switch in the event-ingest handler); "info" is
// then dropped by severityNotifies (felhom.eu/hub/internal/notify/dispatcher.go), which routes only
// warning/error/critical. So a severity outside that set is stored and emailed to NOBODY — neither
// the customer nor the operator leg.
//
// Exported so any caller — and any test in any package — can check the contract against the two
// named hub locations instead of duplicating the literal.
//
// Until v0.215.0 this function emitted "warn", which is not in the set. Every Figyelmeztetés-level
// disk alert the product ever produced was filed as an informational notice and delivered to no one.
func (k DiskAlertKind) Severity() string {
if k == DiskAlertWarn {
return "warning"
}
return "critical"
}
// NotifyDiskHealthDegraded fires a disk_health_degraded hub event. The controller's periodic check
// owns the decision to call this at all (transitions, flap damping, the re-alert cooldown) — never
// first-run, never recovery, never UNKNOWN.
//
// The hub applies its own per-event-type cooldown ON TOP of ours. NOTE: the event type
// "disk_health_degraded" MUST be in the hub's allowedEventTypes (else the hub 400s the POST) — it is.
func (n *Notifier) NotifyDiskHealthDegraded(a DiskAlert) {
critical := a.Kind != DiskAlertWarn
var msg string
switch a.Kind {
case DiskAlertFailSelfReported:
msg = fmt.Sprintf("Lemez állapot romlás: %s — a lemez SMART önellenőrzése hibát jelez. Kérjük, mentse az adatait, és vegye fel velünk a kapcsolatot.", a.Label)
case DiskAlertFailSectors:
msg = fmt.Sprintf("Lemez hiba: %s — a meghajtón %d olvashatatlan szektor van. Mentse az adatait, és keressen meg minket a meghajtó cseréjéhez.", a.Label, a.Sectors)
case DiskAlertFailTemperature:
msg = fmt.Sprintf("Lemez hiba: %s — a meghajtó túlmelegedett (%d °C). Ellenőrizze a gép szellőzését, és keressen meg minket.", a.Label, a.TemperatureC)
case DiskAlertFailWorsened:
msg = fmt.Sprintf("Lemez hiba: %s — a meghajtó állapota tovább romlott, már %d olvashatatlan szektor van. Ha még nem tette meg, mentse az adatait.", a.Label, a.Sectors)
default: // DiskAlertWarn
if len(a.Attributes) > 0 {
msg = fmt.Sprintf("Lemez állapot romlás: %s — romló érték: %s. Javasolt figyelemmel kísérni.", a.Label, strings.Join(a.Attributes, ", "))
} else {
msg = fmt.Sprintf("Lemez állapot romlás: %s — a lemez állapota romlott. Javasolt figyelemmel kísérni.", a.Label)
}
}
n.emit("disk_health_degraded", a.Kind.Severity(), msg,
DiskHealthDetails{Disk: a.Label, Attributes: a.Attributes, Critical: critical})
}
// emit sends an event through the test seam if set, else the real async PushEvent.
//
// Hungarian only: its one producer (the disk-health alerts) still composes a finished sentence.
// It passes no household copy, which is what an un-converted producer looks like.
func (n *Notifier) emit(eventType, severity, message string, details interface{}) {
if n.pushFn != nil {
n.pushFn(eventType, severity, message, "", details)
return
}
n.PushEvent(eventType, severity, message, details)
}
// NotifyAppRemoved sends an app removal event.
func (n *Notifier) NotifyAppRemoved(stackName, displayName string) {
n.pushEventMsg("app_removed", "info", "event.app_removed",
AppDetails{StackName: stackName, DisplayName: displayName}, displayName)
}
// NotifyCrossDriveCompleted sends a cross-drive backup success event.
func (n *Notifier) NotifyCrossDriveCompleted(details CrossDriveDetails) {
n.pushEventMsg("crossdrive_completed", "info", "event.crossdrive_completed", details, details.StackName)
}
// NotifyCrossDriveFailed sends a cross-drive backup failure event.
func (n *Notifier) NotifyCrossDriveFailed(details CrossDriveDetails) {
n.pushEventMsg("crossdrive_failed", "error", "event.crossdrive_failed", details, details.StackName)
}
// NotifyDRStarted sends a disaster recovery start event.
func (n *Notifier) NotifyDRStarted(appCount int) {
n.pushEventMsg("disaster_recovery_started", "warning", "event.disaster_recovery_started", nil, appCount)
}
// NotifyDRCompleted sends a disaster recovery completion event.
func (n *Notifier) NotifyDRCompleted(successCount, failCount int) {
severity := "info"
if failCount > 0 {
severity = "warning"
}
n.pushEventMsg("disaster_recovery_completed", severity, "event.disaster_recovery_completed", nil,
successCount, failCount)
}
// ── Preferences sync ─────────────────────────────────────────────────
type preferencesRequest struct {
CustomerID string `json:"customer_id"`
Email string `json:"email"`
EnabledEvents []string `json:"enabled_events"`
CooldownHours int `json:"cooldown_hours,omitempty"`
}
// SyncPreferences pushes the current notification preferences to the hub.
// Synchronous — returns error for the handler to display to the user.
func (n *Notifier) SyncPreferences(email string, enabledEvents []string, cooldownHours int) error {
if !n.enabled {
return util.MsgError("err.notify.hub_nem_konfiguralt")
}
payload := preferencesRequest{
CustomerID: n.customerID,
Email: email,
EnabledEvents: enabledEvents,
CooldownHours: cooldownHours,
}
jsonData, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("marshal: %w", err)
}
url := n.hubURL + "/api/v1/preferences"
if n.debug {
n.logger.Printf("[DEBUG] SyncPreferences: url=%s email=%s events=%v cooldown=%dh",
url, email, enabledEvents, cooldownHours)
}
req, err := http.NewRequest("POST", url, bytes.NewReader(jsonData))
if err != nil {
return fmt.Errorf("request: %w", err)
}
req.Header.Set("Authorization", "Bearer "+n.apiKey)
req.Header.Set("Content-Type", "application/json")
resp, err := n.httpClient.Do(req)
if err != nil {
return util.MsgError("err.notify.hub_elerhetetlen", err)
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
return util.MsgError("err.notify.hub_error", resp.StatusCode, string(body))
}
if n.debug {
n.logger.Printf("[DEBUG] SyncPreferences: response HTTP %d", resp.StatusCode)
}
n.logger.Printf("[INFO] Notification preferences synced to hub: email=%s, events=%v, cooldown=%dh", email, enabledEvents, cooldownHours)
return nil
}
// ── Test notification ────────────────────────────────────────────────
// SendTest sends a test event for verifying the notification flow (synchronous).
func (n *Notifier) SendTest() error {
if !n.enabled {
return fmt.Errorf("notifications not enabled (hub not configured)")
}
payload := eventRequest{
CustomerID: n.customerID,
EventType: "test",
Severity: "info",
Message: "Teszt értesítés a Felhom rendszerből",
}
jsonData, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("marshal: %w", err)
}
url := n.hubURL + "/api/v1/event"
req, err := http.NewRequest("POST", url, bytes.NewReader(jsonData))
if err != nil {
return fmt.Errorf("request: %w", err)
}
req.Header.Set("Authorization", "Bearer "+n.apiKey)
req.Header.Set("Content-Type", "application/json")
resp, err := n.httpClient.Do(req)
if err != nil {
return fmt.Errorf("send: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
return fmt.Errorf("hub returned %d", resp.StatusCode)
}
return nil
}
// ── Debug event testing ───────────────────────────────────────────────
// PushTestEventSync sends a test event synchronously and returns the Hub HTTP status code.
// Used by the debug page for event testing with configurable type/severity.
func (n *Notifier) PushTestEventSync(eventType, severity, message string) (statusCode int, err error) {
if !n.enabled {
return 0, util.MsgError("err.notify.hub_nem_konfiguralt")
}
payload := eventRequest{
CustomerID: n.customerID,
EventType: eventType,
Severity: severity,
Message: message,
}
jsonData, err := json.Marshal(payload)
if err != nil {
return 0, fmt.Errorf("marshal: %w", err)
}
url := n.hubURL + "/api/v1/event"
req, err := http.NewRequest("POST", url, bytes.NewReader(jsonData))
if err != nil {
return 0, fmt.Errorf("request: %w", err)
}
req.Header.Set("Authorization", "Bearer "+n.apiKey)
req.Header.Set("Content-Type", "application/json")
resp, err := n.httpClient.Do(req)
if err != nil {
n.recordHistory(eventType, severity, message, 0, err.Error())
return 0, fmt.Errorf("send: %w", err)
}
io.Copy(io.Discard, resp.Body)
resp.Body.Close()
if resp.StatusCode >= 400 {
n.recordHistory(eventType, severity, message, resp.StatusCode, fmt.Sprintf("HTTP %d", resp.StatusCode))
return resp.StatusCode, fmt.Errorf("hub returned %d", resp.StatusCode)
}
n.recordHistory(eventType, severity, message, resp.StatusCode, "")
return resp.StatusCode, nil
}
// GetEventHistory returns the last N event history entries (newest first).
func (n *Notifier) GetEventHistory(limit int) []EventHistoryEntry {
n.historyMu.RLock()
defer n.historyMu.RUnlock()
total := n.histPos
if n.histFull {
total = len(n.history)
}
if limit <= 0 || limit > total {
limit = total
}
result := make([]EventHistoryEntry, 0, limit)
for i := 0; i < limit; i++ {
idx := n.histPos - 1 - i
if idx < 0 {
idx += len(n.history)
}
result = append(result, n.history[idx])
}
return result
}
// recordHistory appends an entry to the event history ring buffer.
func (n *Notifier) recordHistory(eventType, severity, message string, hubStatus int, hubError string) {
n.historyMu.Lock()
defer n.historyMu.Unlock()
n.history[n.histPos] = EventHistoryEntry{
Timestamp: time.Now(),
EventType: eventType,
Severity: severity,
Message: message,
HubStatus: hubStatus,
HubError: hubError,
}
n.histPos++
if n.histPos >= len(n.history) {
n.histPos = 0
n.histFull = true
}
}
// ── Backward compatibility ───────────────────────────────────────────
// notifyRequest is the JSON payload for the legacy /api/v1/notify endpoint.
type notifyRequest struct {
CustomerID string `json:"customer_id"`
EventType string `json:"event_type"`
Severity string `json:"severity"`
Message string `json:"message"`
Details string `json:"details,omitempty"`
}
// Notify sends a legacy notification to /api/v1/notify (backward compat).
// Kept for old Hub instances that don't support /api/v1/event yet.
// No local cooldown — Hub handles cooldowns.
func (n *Notifier) Notify(eventType, severity, message, details string) {
if !n.enabled {
return
}
go func() {
payload := notifyRequest{
CustomerID: n.customerID,
EventType: eventType,
Severity: severity,
Message: message,
Details: details,
}
jsonData, err := json.Marshal(payload)
if err != nil {
n.logger.Printf("[ERROR] Failed to marshal notification: %v", err)
return
}
url := n.hubURL + "/api/v1/notify"
req, err := http.NewRequest("POST", url, bytes.NewReader(jsonData))
if err != nil {
return
}
req.Header.Set("Authorization", "Bearer "+n.apiKey)
req.Header.Set("Content-Type", "application/json")
resp, err := n.httpClient.Do(req)
if err != nil {
return
}
io.Copy(io.Discard, resp.Body)
resp.Body.Close()
}()
}
// ── Helpers ──────────────────────────────────────────────────────────
func statusRank(status string) int {
switch status {
case "ok":
return 0
case "warn":
return 1
case "fail":
return 2
default:
return 0
}
}
// NotifyWholeGuestBackupFailed / ...Recovered — R-97a, the WHOLE-GUEST (vzdump) backup tier.
//
// OPERATOR-TIER ONLY, and that is why these are NOT `backup_failed`. `backup_failed` and
// `backup_completed` both carry `customerMessages` entries in the hub AND sit in demo-felhom's live
// `enabled_events`, so reusing them would email the CUSTOMER, in Hungarian, that their backup failed
// — while it is still retrying behind the R-88 breaker. A customer can take no action on a failed
// whole-guest backup; that is the same harm R-97b removes, re-introduced through the front door.
//
// Operator-only is enforced hub-side by `notify.operatorOnlyEvents` (hub >= v0.79.0, R-97c), NOT by
// the absence of a `customerMessages` entry — v0.177.0 claimed the latter and was WRONG: the hub
// falls back to the raw message when the entry is missing, and the only customer gate is
// `prefs.EnabledEvents`, which is configuration. Adding a type to the allowlist does NOT make it
// operator-only; it must go in that register too.
//
// HUB DEPENDENCY: both types MUST be present in the hub's allowedEventTypes or POST /event 400s
// (the recorded allowlist gotcha). Do not deploy this controller ahead of that hub change.
func (n *Notifier) NotifyWholeGuestBackupFailed(tier, message, errMsg string) {
n.PushEvent("whole_guest_backup_failed", "error", message,
WholeGuestBackupDetails{Tier: tier, Error: errMsg})
}
// NotifyBackupTierSkipped — R-518 (v0.243.0). A whole-guest tier was not attempted because its
// storage does not exist on the host. Operator-only hub-side (hub >= v0.114.0 registers the type in
// allowedEventTypes AND operatorOnlyEvents): the household can do nothing about an unprovisioned DR tier.
func (n *Notifier) NotifyBackupTierSkipped(tier, message string) {
n.PushEvent("backup_tier_skipped", "warning", message, WholeGuestBackupDetails{Tier: tier})
}
func (n *Notifier) NotifyWholeGuestBackupRecovered(tier, message string) {
n.PushEvent("whole_guest_backup_recovered", "info", message,
WholeGuestBackupDetails{Tier: tier})
}