From a277b1898167a2842fcc77cbfbab03f7e753cc42 Mon Sep 17 00:00:00 2001 From: kisfenyo Date: Mon, 29 Jun 2026 20:28:02 +0200 Subject: [PATCH] channelhealth: controller->agent channel health-check (periodic probe + classified operator alert) v0.90.0 New internal/channelhealth Checker: ~60s probe via the PRODUCTION memoized client (Server.ProbeAgentChannel, GET /storage), classifies failures (spike Q1 map), debounces transient reasons (N>=2; construction error latches distinctly), seeds first obs, alerts operator+dashboard on transition. Notifier.NotifyAgentChannelDown/Recovered (English, operator-only), AlertManager dashboard banner (Hungarian). No agent/hub change. Spike-proven. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01EPZ4GJ8L5Jqf8UiPwbn1kt --- CHANGELOG.md | 28 +++ controller/cmd/controller/main.go | 34 ++- controller/internal/channelhealth/checker.go | 189 +++++++++++++++ .../internal/channelhealth/checker_test.go | 218 ++++++++++++++++++ controller/internal/notify/notifier.go | 21 +- .../internal/web/agent_disk_handlers.go | 17 ++ controller/internal/web/alerts.go | 30 ++- 7 files changed, 531 insertions(+), 6 deletions(-) create mode 100644 controller/internal/channelhealth/checker.go create mode 100644 controller/internal/channelhealth/checker_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 391f4da..a2d6e84 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,33 @@ ## Changelog +### v0.90.0 — Controller→agent channel health-check (periodic probe + classified operator alert) (2026-06-29) +- **What:** the next self-health slice — a ~60s scheduler job that proves the controller↔agent + local-API channel, classifies failures, and alerts the operator + dashboard on a state change. + Closes the gap the R1 pin-mismatch incident exposed (the channel was only checked once at startup and + only logged). Spike-proven: `felhom.eu/documentation/audits/SPIKE-controller-agent-channel-health-2026-06-29.md`. +- **`internal/channelhealth` (NEW):** a `Checker` over two seams — a `Probe` (the channel call) and a + `Sink` (dashboard + operator notify). Each run classifies the result into `up | down:` by the + spike Q1 error-substring map (`pin_mismatch` / `unauthorized` / `unreachable` / `timeout` / + `misconfigured` / `construction_error` / `unknown`). **Debounce:** transient reasons (refused/timeout) + require **N≥2 consecutive** down-probes before alerting — so the ~1s agent-restart socket gap (spike + Q2) does NOT page; pin/401/DNS/construction alert on the first down observation. First scheduler + observation **seeds** state without notifying (mirrors host_staleness/host_capability). A + construction error (`agentClient()` can't build — latches via `sync.Once`) is surfaced **distinctly**. +- **Probe via the PRODUCTION memoized client (`Server.ProbeAgentChannel`):** GET /storage through + `s.agentClient()` — NOT a fresh `agentapi.New` per probe. Per the spike, the memoized client + self-heals after an agent restart, reflects exactly what the disk UI sees (zero divergence), and + avoids the per-call transport leak the singleton fixed. +- **Operator alert + dashboard:** `Notifier.NotifyAgentChannelDown/Recovered` relay an **English, + operator-only** event (the customer can't act on "the agent re-keyed"; the event type isn't a + customer toggle, same as the host_* events) on a **transition** (up→down, down→up, or reason-change), + with the existing hub 1h cooldown. `AlertManager.SetAgentChannelAlert` shows a short Hungarian + banner whenever the channel is down (state-based, idempotent — a born-down channel shows even though + it's seeded silently). **No customer email.** No agent or hub change (the hub relays the new event + types generically). +- Tests: classifier per-reason, transitions (up→down once, no duplicate, recovered, reason-change + re-alerts), first-obs seed, and the **debounce red-proof** (one refused → no alert; two consecutive → + exactly one). Version `0.89.0 → 0.90.0`. + ### v0.89.0 — App-email: plaintext-only listener (:2526) + split-From mapping (2026-06-29) - **What:** closes the two relay gaps from `FINDING-app-email-rollout-2026-06-29.md` so the opportunistic-STARTTLS clients (cal.com, nextcloud) can use the relay. diff --git a/controller/cmd/controller/main.go b/controller/cmd/controller/main.go index d4bf301..8f58df8 100644 --- a/controller/cmd/controller/main.go +++ b/controller/cmd/controller/main.go @@ -25,6 +25,7 @@ import ( "gitea.dooplex.hu/admin/felhom-controller/internal/assets" "gitea.dooplex.hu/admin/felhom-controller/internal/backup" "gitea.dooplex.hu/admin/felhom-controller/internal/bootstrap" + "gitea.dooplex.hu/admin/felhom-controller/internal/channelhealth" cf "gitea.dooplex.hu/admin/felhom-controller/internal/cloudflare" "gitea.dooplex.hu/admin/felhom-controller/internal/config" "gitea.dooplex.hu/admin/felhom-controller/internal/crypto" @@ -253,9 +254,9 @@ func main() { TLSAddr: cfg.MailRelay.TLSListen, PlainNoTLSAddr: cfg.MailRelay.PlainNoTLSListen, ServiceName: cfg.MailRelay.ShimHost, - Policy: mailrelay.NewPolicy(cfg.MailRelay.FromDomains), - Forwarder: mailrelay.NewHubForwarder(cfg.Hub.URL, cfg.Hub.APIKey), - Logger: logger, + Policy: mailrelay.NewPolicy(cfg.MailRelay.FromDomains), + Forwarder: mailrelay.NewHubForwarder(cfg.Hub.URL, cfg.Hub.APIKey), + Logger: logger, }) if err := mailShim.Apply(sett.AppEmailEnabled()); err != nil { logger.Printf("[ERROR] [mailrelay] could not start app-email shim: %v", err) @@ -717,6 +718,17 @@ func main() { } webServer.SetStartTime(startTime) + // Controller→agent channel health (self-health slice). A ~60s probe of the local-API channel via + // the PRODUCTION memoized client (webServer.ProbeAgentChannel), classified + debounced, alerting the + // operator + dashboard on a transition. Only on a provisioned guest (endpoint set) — mirrors + // probeLocalAPI's guard. Registered after webServer (it owns the memoized client); sched.Every + // launches the job immediately since the scheduler is already started. + if cfg.LocalAPI.Endpoint != "" { + chSink := channelSink{notifier: notifier, alertMgr: alertMgr} + chChecker := channelhealth.New(webServer.ProbeAgentChannel, chSink, logger) + sched.Every("agent-channel-health", 60*time.Second, chChecker.Check) + } + // Wire debug callbacks (only in debug mode) if cfg.Logging.Level == "debug" { dc := &web.DebugCallbacks{} @@ -1308,6 +1320,22 @@ func parseDurationOr(s string, def time.Duration) time.Duration { return d } +// channelSink adapts the notify.Notifier + web.AlertManager to the channelhealth.Sink seam (kept in +// main.go so the channelhealth package imports neither — no cycle). SetDashboard reflects the current +// state each probe (idempotent); NotifyDown/NotifyRecovered fire only on a transition. +type channelSink struct { + notifier *notify.Notifier + alertMgr *web.AlertManager +} + +func (s channelSink) SetDashboard(down bool, _ channelhealth.Reason, msg string) { + s.alertMgr.SetAgentChannelAlert(down, msg) +} +func (s channelSink) NotifyDown(reason channelhealth.Reason, eventType, severity, msg string) { + s.notifier.NotifyAgentChannelDown(string(reason), eventType, severity, msg) +} +func (s channelSink) NotifyRecovered() { s.notifier.NotifyAgentChannelRecovered() } + // probeLocalAPI proves the controller↔agent local-API channel at startup and logs this guest's // mounts (slice 8A). Non-fatal: it only runs when a local-API endpoint is configured, and any // error is logged for diagnosis without affecting the controller's boot. The leaf SHA-256 from diff --git a/controller/internal/channelhealth/checker.go b/controller/internal/channelhealth/checker.go new file mode 100644 index 0000000..ef72f02 --- /dev/null +++ b/controller/internal/channelhealth/checker.go @@ -0,0 +1,189 @@ +// Package channelhealth periodically proves the controller→agent local-API channel and surfaces a +// classified operator alert + dashboard entry on a state change. It closes the gap the R1 +// pin-mismatch incident exposed: the channel was only probed once at startup (probeLocalAPI) and only +// logged, so a mid-life agent re-key went unnoticed until a user hit the dead disk UI. +// +// Design is spike-driven (SPIKE-controller-agent-channel-health-2026-06-29, GO): +// - Probe via the PRODUCTION memoized agent client (reused, not a fresh one): it self-heals after an +// agent restart, reflects exactly what the UI sees, and avoids the per-call transport leak the +// singleton fixed. The construction error (bad fingerprint → agentClient() fails, latching via +// sync.Once) is surfaced distinctly. +// - Classify the failure by the error substring (the spike Q1 map). +// - Debounce transient reasons (connection-refused / timeout): a clean agent restart shows a ~1s +// socket gap (one failed probe) that must NOT page — require N>=2 consecutive before alerting. +// Pin-mismatch / 401 / construction / DNS alert on first observation (they don't self-clear). +// - First scheduler observation SEEDS state without notifying (mirrors host_staleness/host_capability). +package channelhealth + +import ( + "context" + "log" + "strings" + "sync" + "time" +) + +// debounceThreshold is the number of consecutive down-probes a TRANSIENT reason needs before it is +// treated as a real outage (suppresses the ~1s agent-restart blip from the spike's Q2). +const debounceThreshold = 2 + +// probeTimeout mirrors probeLocalAPI's 15s budget for the GET /storage call. +const probeTimeout = 15 * time.Second + +// Reason is the classified channel-down cause. +type Reason string + +const ( + ReasonPinMismatch Reason = "pin_mismatch" + ReasonUnauthorized Reason = "unauthorized" + ReasonUnreachable Reason = "unreachable" + ReasonTimeout Reason = "timeout" + ReasonMisconfigured Reason = "misconfigured" // endpoint unresolvable (DNS) + ReasonConstruction Reason = "construction_error" // agentClient() build error — LATCHES until restart + ReasonUnknown Reason = "unknown" +) + +// Probe runs one channel check via the production memoized client. It returns whether the failure was +// a CONSTRUCTION error (agentClient() couldn't even build — a latching config fault) and the error +// (nil = channel up). Production wires this to Server.ProbeAgentChannel; tests inject a fake. +type Probe func(ctx context.Context) (constructionErr bool, err error) + +// Sink receives the checker's outputs. SetDashboard reflects the CURRENT state every check (idempotent +// — a born-down channel shows immediately, with no notification). NotifyDown/NotifyRecovered fire ONLY +// on a transition (the operator relay, with the hub-side cooldown). Implemented in main.go over the +// notify.Notifier + web.AlertManager (kept out of this package to avoid an import cycle). +type Sink interface { + SetDashboard(down bool, reason Reason, hungarianMsg string) + NotifyDown(reason Reason, eventType, severity, englishMsg string) + NotifyRecovered() +} + +type classification struct { + reason Reason + eventType string + severity string // "error" (critical) | "warning" (transient) + english string // operator alert message (relayed to the hub) + hungarian string // short dashboard line + debounce bool // transient → require N>=2 consecutive +} + +// classify maps a probe result to a classification. constructionErr is authoritative for the latching +// config fault; otherwise the error substring decides (spike Q1 map). Match substrings, not exact +// strings (the wrapped Get "..." prefix varies). +func classify(constructionErr bool, err error) classification { + if constructionErr { + return classification{ReasonConstruction, "agent_channel_construction_error", "error", + "Controller→agent channel down: agent client misconfigured — local_api.fingerprint is not a valid SHA-256, the controller cannot build the client (config fix + controller restart needed).", + "A tárolókezelő ügynök kapcsolata hibásan beállítva.", false} + } + s := err.Error() + switch { + case strings.Contains(s, "TLS pin mismatch"): + return classification{ReasonPinMismatch, "agent_channel_pin_mismatch", "error", + "Controller→agent channel down: TLS pin mismatch — the agent's leaf cert no longer matches the controller's bootstrap fingerprint (re-pin / re-bootstrap).", + "A tárolókezelő ügynök tanúsítványa megváltozott.", false} + case strings.Contains(s, "HTTP 401"): + return classification{ReasonUnauthorized, "agent_channel_unauthorized", "error", + "Controller→agent channel down: agent rejected the controller token (HTTP 401) — token stale/rotated (re-bootstrap).", + "A tárolókezelő ügynök elutasította a hozzáférést.", false} + case strings.Contains(s, "no such host") || strings.Contains(s, "lookup "): + return classification{ReasonMisconfigured, "agent_channel_misconfigured", "error", + "Controller→agent channel down: agent endpoint unresolvable — local_api.endpoint misconfigured.", + "A tárolókezelő ügynök címe nem feloldható.", false} + case strings.Contains(s, "connection refused"): + return classification{ReasonUnreachable, "agent_channel_unreachable", "warning", + "Controller→agent channel down: agent unreachable (connection refused) — felhom-agent down or :8443 closed.", + "A tárolókezelő ügynök nem elérhető.", true} + case strings.Contains(s, "context deadline exceeded") || strings.Contains(s, "i/o timeout") || strings.Contains(s, "Client.Timeout"): + return classification{ReasonTimeout, "agent_channel_timeout", "warning", + "Controller→agent channel down: timeout — host unreachable / packets dropped.", + "A tárolókezelő ügynök nem válaszol.", true} + default: + return classification{ReasonUnknown, "agent_channel_unknown", "warning", + "Controller→agent channel down: " + s, + "A tárolókezelő ügynök nem elérhető.", true} + } +} + +// Checker holds the in-memory channel state. No persistence — the state is re-derived each run +// (mirrors the AlertManager's state-based model). Safe for the single scheduler caller; the mutex +// guards against an overlapping run. +type Checker struct { + probe Probe + sink Sink + logger *log.Logger + + mu sync.Mutex + state string // "" (unseeded) | "up" | "down:" + consecutiveDown int +} + +// New builds a checker over the probe + sink seams. +func New(probe Probe, sink Sink, logger *log.Logger) *Checker { + return &Checker{probe: probe, sink: sink, logger: logger} +} + +// Check runs one probe cycle: classify, debounce, reflect on the dashboard, and notify on a real +// transition. Best-effort + idempotent — safe to call on a timer. Never returns an error to the +// scheduler (a probe failure IS the signal, not a job failure). +func (c *Checker) Check(ctx context.Context) error { + pctx, cancel := context.WithTimeout(ctx, probeTimeout) + constructionErr, perr := c.probe(pctx) + cancel() + + c.mu.Lock() + defer c.mu.Unlock() + + // Channel UP. + if perr == nil { + c.consecutiveDown = 0 + c.sink.SetDashboard(false, "", "") + prev := c.state + c.state = "up" + if prev == "" || prev == "up" { + return nil // seed, or steady-up → no notify + } + c.logger.Printf("[INFO] [channel] agent channel recovered (was %s)", prev) + c.sink.NotifyRecovered() + return nil + } + + // Channel DOWN — classify + debounce transient reasons. + cls := classify(constructionErr, perr) + c.consecutiveDown++ + if cls.debounce && c.consecutiveDown < debounceThreshold { + // A transient blip (e.g. the ~1s agent-restart socket gap). Hold the previous state — do NOT + // flip the dashboard or notify. If we've never seen anything yet, assume up until confirmed. + if c.state == "" { + c.state = "up" + } + c.logger.Printf("[DEBUG] [channel] transient down (%s, %d/%d) — suppressed pending confirmation: %v", + cls.reason, c.consecutiveDown, debounceThreshold, perr) + return nil + } + + newState := "down:" + string(cls.reason) + c.sink.SetDashboard(true, cls.reason, cls.hungarian) // dashboard reflects current state always + prev := c.state + c.state = newState + if prev == "" { + c.logger.Printf("[INFO] [channel] agent channel down at startup (%s) — seeded, no alert: %v", cls.reason, perr) + return nil // first observation seeds, no alert (dashboard already set above) + } + if prev == newState { + return nil // steady down, same reason → no duplicate notify (the dashboard stays set) + } + c.logger.Printf("[WARN] [channel] agent channel DOWN (%s→%s): %v", prev, newState, perr) + c.sink.NotifyDown(cls.reason, cls.eventType, cls.severity, cls.english) + return nil +} + +// State returns the current channel state (for tests/diagnostics). +func (c *Checker) State() string { + c.mu.Lock() + defer c.mu.Unlock() + if c.state == "" { + return "unseeded" + } + return c.state +} diff --git a/controller/internal/channelhealth/checker_test.go b/controller/internal/channelhealth/checker_test.go new file mode 100644 index 0000000..2f0c450 --- /dev/null +++ b/controller/internal/channelhealth/checker_test.go @@ -0,0 +1,218 @@ +package channelhealth + +import ( + "context" + "errors" + "io" + "log" + "testing" +) + +type downCall struct { + reason Reason + eventType string + severity string +} + +type fakeSink struct { + downs []downCall + recovered int + dashDown bool + dashReason Reason +} + +func (f *fakeSink) SetDashboard(down bool, reason Reason, _ string) { + f.dashDown = down + f.dashReason = reason +} +func (f *fakeSink) NotifyDown(reason Reason, eventType, severity, _ string) { + f.downs = append(f.downs, downCall{reason, eventType, severity}) +} +func (f *fakeSink) NotifyRecovered() { f.recovered++ } + +// scriptedProbe returns the next (constructionErr, err) on each call; a fresh agentapi.New is NEVER +// built here — the checker only calls this seam (proving it reuses the production client). +type scriptedProbe struct { + steps []struct { + cons bool + err error + } + i int + calls int +} + +func (p *scriptedProbe) fn(_ context.Context) (bool, error) { + p.calls++ + s := p.steps[p.i] + if p.i < len(p.steps)-1 { + p.i++ + } + return s.cons, s.err +} + +func newChecker(t *testing.T, sink Sink) *Checker { + t.Helper() + return New(nil, sink, log.New(io.Discard, "", 0)) +} + +func step(cons bool, err error) struct { + cons bool + err error +} { + return struct { + cons bool + err error + }{cons, err} +} + +// run feeds the checker a sequence of probe results. +func run(c *Checker, p *scriptedProbe) { + c.probe = p.fn + for range p.steps { + _ = c.Check(context.Background()) + } +} + +// §8 classification: each error → correct reason/event/severity. Seed up first so the down is a +// transition (non-debounce reasons fire on the first down observation). +func TestClassify_PerReason(t *testing.T) { + cases := []struct { + name string + cons bool + err error + reason Reason + eventType string + severity string + }{ + {"pin", false, errors.New(`agentapi: GET /storage: Get "https://x": agentapi: TLS pin mismatch: ...`), ReasonPinMismatch, "agent_channel_pin_mismatch", "error"}, + {"401", false, errors.New("agentapi: GET /storage: HTTP 401"), ReasonUnauthorized, "agent_channel_unauthorized", "error"}, + {"dns", false, errors.New(`agentapi: GET /storage: Get "https://x": dial tcp: lookup x: no such host`), ReasonMisconfigured, "agent_channel_misconfigured", "error"}, + {"construction", true, errors.New("agentapi: fingerprint must be a SHA-256 (64 hex chars), got 10"), ReasonConstruction, "agent_channel_construction_error", "error"}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + sink := &fakeSink{} + c := newChecker(t, sink) + p := &scriptedProbe{steps: []struct { + cons bool + err error + }{step(false, nil), step(tc.cons, tc.err)}} // seed up → down + run(c, p) + if len(sink.downs) != 1 { + t.Fatalf("want 1 down alert, got %d", len(sink.downs)) + } + d := sink.downs[0] + if d.reason != tc.reason || d.eventType != tc.eventType || d.severity != tc.severity { + t.Fatalf("got %+v, want reason=%s event=%s sev=%s", d, tc.reason, tc.eventType, tc.severity) + } + if !sink.dashDown { + t.Errorf("dashboard should be down") + } + }) + } +} + +// §7-C debounce RED-PROOF: a single connection-refused (the ~1s agent-restart blip) must NOT alert; +// TWO consecutive must alert exactly once. A no-debounce impl fires on the single blip → fails this. +func TestDebounce_TransientBlipSuppressed(t *testing.T) { + refused := errors.New(`agentapi: GET /storage: Get "https://x": dial tcp: connect: connection refused`) + + // One refused sandwiched by up: NO alert. + sink := &fakeSink{} + c := newChecker(t, sink) + run(c, &scriptedProbe{steps: []struct { + cons bool + err error + }{step(false, nil), step(false, refused), step(false, nil)}}) + if len(sink.downs) != 0 { + t.Fatalf("single transient blip must NOT alert, got %d", len(sink.downs)) + } + if sink.dashDown { + t.Errorf("single blip must not flip the dashboard down") + } + + // Two consecutive refused: exactly ONE alert. + sink2 := &fakeSink{} + c2 := newChecker(t, sink2) + run(c2, &scriptedProbe{steps: []struct { + cons bool + err error + }{step(false, nil), step(false, refused), step(false, refused)}}) + if len(sink2.downs) != 1 || sink2.downs[0].reason != ReasonUnreachable { + t.Fatalf("two consecutive refused → want 1 unreachable alert, got %+v", sink2.downs) + } +} + +// §7-A/B: up→down alerts once (not per cycle); §7: down→up recovers. +func TestTransitions_NoDuplicate_AndRecovery(t *testing.T) { + pin := errors.New("agentapi: GET /storage: ...: agentapi: TLS pin mismatch: ...") + sink := &fakeSink{} + c := newChecker(t, sink) + // seed up, then 3× down (same reason), then up. + run(c, &scriptedProbe{steps: []struct { + cons bool + err error + }{step(false, nil), step(false, pin), step(false, pin), step(false, pin), step(false, nil)}}) + if len(sink.downs) != 1 { + t.Fatalf("steady down must alert once, got %d", len(sink.downs)) + } + if sink.recovered != 1 { + t.Fatalf("want 1 recovered, got %d", sink.recovered) + } + if sink.dashDown { + t.Errorf("dashboard should be cleared after recovery") + } +} + +// §7-E: a reason change (unreachable→pin_mismatch) re-alerts. +func TestReasonChange_ReAlerts(t *testing.T) { + refused := errors.New("...: connect: connection refused") + pin := errors.New("...: agentapi: TLS pin mismatch: ...") + sink := &fakeSink{} + c := newChecker(t, sink) + // seed up; 2× refused (confirm unreachable → alert 1); then pin (reason change → alert 2). + run(c, &scriptedProbe{steps: []struct { + cons bool + err error + }{step(false, nil), step(false, refused), step(false, refused), step(false, pin)}}) + if len(sink.downs) != 2 { + t.Fatalf("reason change should re-alert (want 2), got %d: %+v", len(sink.downs), sink.downs) + } + if sink.downs[0].reason != ReasonUnreachable || sink.downs[1].reason != ReasonPinMismatch { + t.Fatalf("want unreachable then pin_mismatch, got %+v", sink.downs) + } +} + +// §9.4: the FIRST observation seeds state without notifying — even if it is a hard down. +func TestFirstObservation_SeedsNoAlert(t *testing.T) { + pin := errors.New("...: agentapi: TLS pin mismatch: ...") + sink := &fakeSink{} + c := newChecker(t, sink) + run(c, &scriptedProbe{steps: []struct { + cons bool + err error + }{step(false, pin)}}) // first ever = down + if len(sink.downs) != 0 { + t.Fatalf("first observation must not notify, got %d", len(sink.downs)) + } + if !sink.dashDown { + t.Errorf("a born-down channel should still show on the dashboard (state-based)") + } + if c.State() != "down:pin_mismatch" { + t.Errorf("state = %s, want down:pin_mismatch", c.State()) + } +} + +// The checker only uses the injected probe seam (never a fresh agentapi.New). +func TestProbe_SeamOnly(t *testing.T) { + sink := &fakeSink{} + c := newChecker(t, sink) + p := &scriptedProbe{steps: []struct { + cons bool + err error + }{step(false, nil)}} + run(c, p) + if p.calls != 1 { + t.Fatalf("checker should call the probe seam exactly once, got %d", p.calls) + } +} diff --git a/controller/internal/notify/notifier.go b/controller/internal/notify/notifier.go index 61742c8..e828b09 100644 --- a/controller/internal/notify/notifier.go +++ b/controller/internal/notify/notifier.go @@ -335,6 +335,26 @@ func (n *Notifier) NotifyStorageReconnected(label string) { fmt.Sprintf("Meghajtó újra csatlakoztatva: %s", label), StorageDetails{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) +} + // NotifyAppDeployed sends an app deployment event. func (n *Notifier) NotifyAppDeployed(stackName, displayName string) { n.PushEvent("app_deployed", "info", @@ -630,4 +650,3 @@ func statusRank(status string) int { return 0 } } - diff --git a/controller/internal/web/agent_disk_handlers.go b/controller/internal/web/agent_disk_handlers.go index b2f0133..9ef7fcf 100644 --- a/controller/internal/web/agent_disk_handlers.go +++ b/controller/internal/web/agent_disk_handlers.go @@ -1,6 +1,7 @@ package web import ( + "context" "encoding/json" "errors" "net/http" @@ -59,6 +60,22 @@ func (s *Server) agentClient() (*agentapi.Client, error) { return s.agentCli, s.agentCliErr } +// ProbeAgentChannel runs one controller→agent channel health probe using the PRODUCTION memoized +// client (NOT a fresh one — spike SPIKE-controller-agent-channel-health-2026-06-29: it self-heals, +// reflects exactly what the disk UI sees, and avoids the per-call transport leak the singleton fixed). +// It returns whether the failure was a CONSTRUCTION error (agentClient() couldn't build — a latching +// config fault, distinct from a runtime channel failure) and the error (nil = channel up). The probe +// is GET /storage (cheap, read-only — the same call probeLocalAPI uses at startup). This is the +// channelhealth.Probe seam. +func (s *Server) ProbeAgentChannel(ctx context.Context) (constructionErr bool, err error) { + client, cerr := s.agentClient() + if cerr != nil { + return true, cerr + } + _, serr := client.Storage(ctx) + return false, serr +} + // writeDiskJSON writes the standard {ok,data,error} envelope used by the disk API. func writeDiskJSON(w http.ResponseWriter, status int, ok bool, errMsg string, data interface{}) { w.Header().Set("Content-Type", "application/json") diff --git a/controller/internal/web/alerts.go b/controller/internal/web/alerts.go index 1ecef05..cfcbf43 100644 --- a/controller/internal/web/alerts.go +++ b/controller/internal/web/alerts.go @@ -33,6 +33,10 @@ type AlertManager struct { alerts []Alert logger *log.Logger hubPushStatusFn func() HubPushStatusData + // agentChannelAlert is set/cleared by the channel-health checker (out-of-band from Refresh, which + // is health-report-driven). nil = channel up. It is prepended in GetAlerts so a dead controller→ + // agent link (disk/storage UI broken) is always visible regardless of the health-report cycle. + agentChannelAlert *Alert } // NewAlertManager creates a new AlertManager. @@ -49,6 +53,25 @@ func (am *AlertManager) SetHubPushStatus(fn func() HubPushStatusData) { am.mu.Unlock() } +// SetAgentChannelAlert sets (down=true) or clears (down=false) the controller→agent channel-down +// dashboard banner. Called by the channel-health checker each probe (idempotent). msg is the short +// Hungarian display line. +func (am *AlertManager) SetAgentChannelAlert(down bool, msg string) { + am.mu.Lock() + defer am.mu.Unlock() + if !down { + am.agentChannelAlert = nil + return + } + am.agentChannelAlert = &Alert{ + ID: "agent-channel-down", + Level: "error", + Message: msg, + Link: "/settings", + LinkText: "Beállítások", + } +} + // Refresh regenerates alerts from the latest health check report and config state. // Called after each health check cycle (every 5 minutes) and on storage state changes. func (am *AlertManager) Refresh(report *monitor.HealthReport, cfg *config.Config, backupMgr *backup.Manager, updateAvailable bool, latestVersion string, storagePaths ...[]settings.StoragePath) { @@ -159,7 +182,7 @@ func (am *AlertManager) GetAlerts(excludeIDs ...string) []Alert { am.mu.RLock() defer am.mu.RUnlock() - if len(am.alerts) == 0 { + if len(am.alerts) == 0 && am.agentChannelAlert == nil { return nil } @@ -169,6 +192,10 @@ func (am *AlertManager) GetAlerts(excludeIDs ...string) []Alert { } var result []Alert + // Channel-down is prepended (highest priority — the agent link being dead breaks disk/storage UI). + if am.agentChannelAlert != nil && !exclude[am.agentChannelAlert.ID] { + result = append(result, *am.agentChannelAlert) + } for _, a := range am.alerts { if exclude[a.ID] { continue @@ -239,4 +266,3 @@ func countLevel(alerts []Alert, level string) int { } return n } -