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) } // The household half is named, not printed. An operator asking "why did my English // customer get a Hungarian mail?" needs to know whether the BOX attached a second // sentence at all — that is the fork in the diagnosis, and without it the answer is // invisible on both sides of the wire (the hub logs only the Hungarian, and the // payload is inside TLS). The sentence itself is not logged: it is the same sentence // twice, and one copy is already on the line. household := "hu-only" if messageCustomer != "" { household = "+household(" + n.boxLang() + ")" } n.logger.Printf("[INFO] Event pushed: %s (%s) [%s] — %s", eventType, severity, household, 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}) }