From 137828a71376de2cd844d468253be5027d6a5a4f Mon Sep 17 00:00:00 2001 From: Stan Misiurev Date: Thu, 6 Aug 2026 11:39:52 -0400 Subject: [PATCH 1/3] feat(monitoring): add channel health monitoring for OpenClaw instances Control plane now polls each running instance's OpenClaw gateway (channels.status WS RPC over the existing Gateway tunnel) to detect channels that are disconnected, wedged (stale sockets), or not running, catching the known failure mode where the process is alive but a Slack Socket Mode / Telegram poller connection is silently dead. - internal/channelhealth: background monitor (60s interval, bounded concurrency), health evaluation with 30m stale threshold, persisted per-channel status + in-memory snapshot - GET /api/v1/instances/{id}/channels/health + channel_health summary on instance responses - Channel Health panel on Agent detail Settings tab and warning indicator on the agent list for unhealthy/unreachable instances - CLAWORC_CHANNEL_HEALTH_ENABLED / CLAWORC_CHANNEL_HEALTH_INTERVAL - docs/channel-monitoring.md spec + end-user page Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01CfW5Dn8yEsULgbKzmcjbZ7 --- .../src/app/pages/AgentDetailPage.tsx | 4 + .../frontend/src/common/api/channels.ts | 7 + .../src/common/components/AgentCard.tsx | 6 +- .../src/common/components/AgentRow.tsx | 20 +- .../common/components/ChannelHealthPanel.tsx | 191 +++++++++ .../src/common/hooks/useChannelHealth.ts | 13 + .../frontend/src/common/types/channel.ts | 47 +++ .../frontend/src/common/types/instance.ts | 4 + .../internal/channelhealth/channelhealth.go | 275 +++++++++++++ .../channelhealth/channelhealth_test.go | 283 ++++++++++++++ .../internal/channelhealth/monitor.go | 363 ++++++++++++++++++ control-plane/internal/config/config.go | 7 + .../migrations/migration_00001_baseline.go | 1 + control-plane/internal/database/models.go | 53 +-- .../internal/database/models/models.go | 23 ++ control-plane/internal/handlers/channels.go | 127 ++++++ control-plane/internal/handlers/instances.go | 36 ++ control-plane/main.go | 12 + docs/README.md | 1 + docs/channel-monitoring.md | 145 +++++++ website_docs/channel-health.mdx | 75 ++++ 21 files changed, 1658 insertions(+), 35 deletions(-) create mode 100644 control-plane/frontend/src/common/api/channels.ts create mode 100644 control-plane/frontend/src/common/components/ChannelHealthPanel.tsx create mode 100644 control-plane/frontend/src/common/hooks/useChannelHealth.ts create mode 100644 control-plane/frontend/src/common/types/channel.ts create mode 100644 control-plane/internal/channelhealth/channelhealth.go create mode 100644 control-plane/internal/channelhealth/channelhealth_test.go create mode 100644 control-plane/internal/channelhealth/monitor.go create mode 100644 control-plane/internal/handlers/channels.go create mode 100644 docs/channel-monitoring.md create mode 100644 website_docs/channel-health.mdx diff --git a/control-plane/frontend/src/app/pages/AgentDetailPage.tsx b/control-plane/frontend/src/app/pages/AgentDetailPage.tsx index 9ea6bbfd..5c30ad32 100644 --- a/control-plane/frontend/src/app/pages/AgentDetailPage.tsx +++ b/control-plane/frontend/src/app/pages/AgentDetailPage.tsx @@ -53,6 +53,7 @@ import AffinityEditor from "@common/components/AffinityEditor"; import PortsEditor from "@common/components/PortsEditor"; import { useHealth } from "@common/hooks/useHealth"; import WebhookSection from "@common/components/WebhookSection"; +import ChannelHealthPanel from "@common/components/ChannelHealthPanel"; import LegacyBrowserBanner from "@common/components/LegacyBrowserBanner"; import AppToast from "@common/components/AppToast"; import { infoToast } from "@common/utils/toast"; @@ -1153,6 +1154,9 @@ export default function AgentDetailPage() { {/* Webhook (per-instance) — admins and team managers */} + {/* Chat channel health */} + + {/* SSH Connection Status */} { + const { data } = await client.get(`/instances/${instanceId}/channels/health`); + return data; +} diff --git a/control-plane/frontend/src/common/components/AgentCard.tsx b/control-plane/frontend/src/common/components/AgentCard.tsx index 75816919..611fa3d3 100644 --- a/control-plane/frontend/src/common/components/AgentCard.tsx +++ b/control-plane/frontend/src/common/components/AgentCard.tsx @@ -3,6 +3,7 @@ import { formatDistanceToNow } from "date-fns"; import { GripVertical } from "lucide-react"; import StatusBadge from "./StatusBadge"; import ActionButtons from "./ActionButtons"; +import { ChannelHealthIndicator } from "./ChannelHealthPanel"; import { useSSHStatus } from "@common/hooks/useSSHStatus"; import { buildSSHTooltip } from "@common/utils/sshTooltip"; import type { Instance } from "@common/types/instance"; @@ -61,7 +62,10 @@ export default function AgentCard({ {instance.display_name} - +
+ + +
{createdAt}
diff --git a/control-plane/frontend/src/common/components/AgentRow.tsx b/control-plane/frontend/src/common/components/AgentRow.tsx index bcfd36c1..a02d5387 100644 --- a/control-plane/frontend/src/common/components/AgentRow.tsx +++ b/control-plane/frontend/src/common/components/AgentRow.tsx @@ -3,6 +3,7 @@ import { formatDistanceToNow } from "date-fns"; import { GripVertical } from "lucide-react"; import StatusBadge from "./StatusBadge"; import ActionButtons from "./ActionButtons"; +import { ChannelHealthIndicator } from "./ChannelHealthPanel"; import { useSSHStatus } from "@common/hooks/useSSHStatus"; import { buildSSHTooltip } from "@common/utils/sshTooltip"; import type { Instance } from "@common/types/instance"; @@ -58,14 +59,17 @@ export default function AgentRow({ - +
+ + +
{createdAt} diff --git a/control-plane/frontend/src/common/components/ChannelHealthPanel.tsx b/control-plane/frontend/src/common/components/ChannelHealthPanel.tsx new file mode 100644 index 00000000..b679de0d --- /dev/null +++ b/control-plane/frontend/src/common/components/ChannelHealthPanel.tsx @@ -0,0 +1,191 @@ +import { RefreshCw } from "lucide-react"; +import { formatDistanceToNow } from "date-fns"; +import { useChannelHealth } from "@common/hooks/useChannelHealth"; +import type { ChannelAccountHealth, ChannelHealthSummary } from "@common/types/channel"; +import type { Instance } from "@common/types/instance"; + +const overallStyles: Record = { + healthy: "bg-green-100 text-green-800", + degraded: "bg-yellow-100 text-yellow-800", + unhealthy: "bg-red-100 text-red-800", + unreachable: "bg-red-100 text-red-800", + no_channels: "bg-gray-100 text-gray-800", + unknown: "bg-gray-100 text-gray-800", +}; + +const channelStatusStyles: Record = { + healthy: "bg-green-100 text-green-800", + stale: "bg-yellow-100 text-yellow-800", + disconnected: "bg-red-100 text-red-800", + not_running: "bg-red-100 text-red-800", + disabled: "bg-gray-100 text-gray-800", + unknown: "bg-gray-100 text-gray-800", +}; + +function statusLabel(status: string): string { + return status.replace(/_/g, " "); +} + +function capitalize(s: string): string { + return s ? s.charAt(0).toUpperCase() + s.slice(1) : s; +} + +function relativeTime(ts: string | null): string | null { + if (!ts) return null; + const d = new Date(ts); + if (isNaN(d.getTime())) return null; + return formatDistanceToNow(d, { addSuffix: true }); +} + +function ChannelRow({ ch }: { ch: ChannelAccountHealth }) { + const badgeStyle = channelStatusStyles[ch.status] ?? "bg-gray-100 text-gray-800"; + const lastEvent = relativeTime(ch.last_event_at); + return ( +
+
+ + {capitalize(ch.channel)} + {ch.account_id && ch.account_id !== "default" && ( + ({ch.account_id}) + )} + + + {statusLabel(ch.status)} + + {ch.mode && {ch.mode}} + + {lastEvent ? `last event ${lastEvent}` : "no events yet"} + +
+ {(ch.last_error || ch.reconnect_attempts > 0) && ( +
+ {ch.last_error && {ch.last_error}} + {ch.reconnect_attempts > 0 && ( + + {ch.reconnect_attempts} reconnect attempt{ch.reconnect_attempts === 1 ? "" : "s"} + + )} +
+ )} +
+ ); +} + +export default function ChannelHealthPanel({ instanceId }: { instanceId: number }) { + const health = useChannelHealth(instanceId); + + if (health.isLoading && !health.data) { + return ( +
+
Loading channel health...
+
+ ); + } + + if (health.isError && !health.data) { + return ( +
+
+
Failed to load channel health.
+ +
+
+ ); + } + + if (!health.data) return null; + + const data = health.data; + + // Monitoring is turned off server-side — hide the panel entirely. + if (data.overall === "disabled") return null; + + const overallStyle = overallStyles[data.overall] ?? "bg-gray-100 text-gray-800"; + const checkedAt = relativeTime(data.checked_at); + + let body; + if (data.overall === "no_channels") { + body =

No channels configured

; + } else if (data.overall === "unreachable") { + body = ( +

+ Gateway unreachable — the OpenClaw process may be down +

+ ); + } else if (data.overall === "unknown") { + body =

Waiting for first health check…

; + } else { + body = ( +
+ {data.channels.map((ch) => ( + + ))} +
+ ); + } + + return ( +
+
+
+

Channel Health

+ + {statusLabel(data.overall)} + +
+
+ {checkedAt && checked {checkedAt}} + +
+
+ {body} +
+ ); +} + +/** + * Compact warning indicator for agent list rows/cards. Renders only when the + * instance's channel health summary is "unhealthy" or "unreachable". + */ +export function ChannelHealthIndicator({ instance }: { instance: Instance }) { + const summary: ChannelHealthSummary | null | undefined = instance.channel_health; + if (!summary) return null; + if (summary.overall !== "unhealthy" && summary.overall !== "unreachable") return null; + + const tooltip = + summary.overall === "unreachable" + ? "Gateway unreachable — the OpenClaw process may be down" + : `${summary.unhealthy_count} channel${summary.unhealthy_count === 1 ? "" : "s"} not responding`; + + return ( + + + + Channels + + + {tooltip} + + + ); +} diff --git a/control-plane/frontend/src/common/hooks/useChannelHealth.ts b/control-plane/frontend/src/common/hooks/useChannelHealth.ts new file mode 100644 index 00000000..1e3e28c5 --- /dev/null +++ b/control-plane/frontend/src/common/hooks/useChannelHealth.ts @@ -0,0 +1,13 @@ +import { useQuery } from "@tanstack/react-query"; +import { getChannelHealth } from "@common/api/channels"; + +export function useChannelHealth(instanceId: number | undefined) { + return useQuery({ + queryKey: ["instance-channel-health", instanceId], + queryFn: () => getChannelHealth(instanceId!), + enabled: !!instanceId, + refetchInterval: 15000, + refetchIntervalInBackground: false, + retry: false, + }); +} diff --git a/control-plane/frontend/src/common/types/channel.ts b/control-plane/frontend/src/common/types/channel.ts new file mode 100644 index 00000000..68e9c65e --- /dev/null +++ b/control-plane/frontend/src/common/types/channel.ts @@ -0,0 +1,47 @@ +export type ChannelOverallStatus = + | "healthy" + | "degraded" + | "unhealthy" + | "unreachable" + | "no_channels" + | "unknown" + | "disabled"; + +export type ChannelAccountStatus = + | "healthy" + | "disconnected" + | "not_running" + | "stale" + | "disabled" + | "unknown"; + +export interface ChannelAccountHealth { + channel: string; + account_id: string; + status: ChannelAccountStatus; + enabled: boolean; + running: boolean; + connected: boolean; + mode: string; + last_event_at: string | null; + last_inbound_at: string | null; + last_outbound_at: string | null; + last_error: string; + reconnect_attempts: number; + checked_at: string | null; +} + +export interface ChannelHealth { + instance_id: number; + overall: ChannelOverallStatus; + gateway_reachable: boolean; + checked_at: string | null; + channels: ChannelAccountHealth[]; +} + +/** Compact summary embedded in Instance list/detail responses. */ +export interface ChannelHealthSummary { + overall: ChannelOverallStatus; + unhealthy_count: number; + checked_at: string | null; +} diff --git a/control-plane/frontend/src/common/types/instance.ts b/control-plane/frontend/src/common/types/instance.ts index 0ec889ea..c732a5d3 100644 --- a/control-plane/frontend/src/common/types/instance.ts +++ b/control-plane/frontend/src/common/types/instance.ts @@ -1,3 +1,5 @@ +import type { ChannelHealthSummary } from "./channel"; + export interface InstanceModels { effective: string[]; disabled_defaults: string[]; @@ -58,6 +60,8 @@ export interface Instance { affinity: string; service_account_annotations: Record; ports: PortSpec[]; + /** Compact chat-channel health summary (absent when monitoring hasn't run or is disabled). */ + channel_health?: ChannelHealthSummary | null; } export interface PortSpec { diff --git a/control-plane/internal/channelhealth/channelhealth.go b/control-plane/internal/channelhealth/channelhealth.go new file mode 100644 index 00000000..12f8796e --- /dev/null +++ b/control-plane/internal/channelhealth/channelhealth.go @@ -0,0 +1,275 @@ +// Package channelhealth implements a background monitor that polls each +// running instance's OpenClaw gateway for per-channel/per-account runtime +// state (via the channels.status RPC), evaluates health, persists the +// latest status to the database, and keeps an in-memory snapshot for +// cheap reads by HTTP handlers. +package channelhealth + +import ( + "bytes" + "encoding/json" + "fmt" + "sort" + "time" +) + +// Per-account health statuses. +const ( + StatusHealthy = "healthy" + StatusDisconnected = "disconnected" + StatusNotRunning = "not_running" + StatusStale = "stale" + StatusDisabled = "disabled" + StatusUnknown = "unknown" +) + +// Instance-level overall statuses. +const ( + OverallHealthy = "healthy" + OverallDegraded = "degraded" + OverallUnhealthy = "unhealthy" + OverallUnreachable = "unreachable" + OverallNoChannels = "no_channels" + OverallUnknown = "unknown" +) + +// StaleThreshold is how long a connected persistent-socket channel may go +// without any event before it is considered stale. +const StaleThreshold = 30 * time.Minute + +// AccountState is the parsed runtime state of one channel account as +// reported by the gateway's channels.status RPC. +type AccountState struct { + AccountID string + Enabled bool + Configured bool + Running bool + Connected bool + Mode string + LastEventAt *time.Time + LastInboundAt *time.Time + LastOutboundAt *time.Time + LastError string + ReconnectAttempts int +} + +// ChannelState is the evaluated health of one channel account, ready for +// persistence and for serving to the frontend. +type ChannelState struct { + Channel string + AccountID string + Status string + Enabled bool + Running bool + Connected bool + Mode string + LastEventAt *time.Time + LastInboundAt *time.Time + LastOutboundAt *time.Time + LastError string + ReconnectAttempts int + CheckedAt time.Time +} + +// Snapshot is the latest known channel health for one instance. +type Snapshot struct { + InstanceID uint + Overall string + GatewayReachable bool + CheckedAt time.Time + Channels []ChannelState +} + +// --- channels.status payload parsing ----------------------------------- + +// wirePayload mirrors the channels.status response payload. Every field is +// optional/nullable on the wire, so everything is a pointer or raw JSON. +type wirePayload struct { + ChannelAccounts map[string][]wireAccount `json:"channelAccounts"` +} + +type wireAccount struct { + AccountID *string `json:"accountId"` + Enabled *bool `json:"enabled"` + Configured *bool `json:"configured"` + Running *bool `json:"running"` + Connected *bool `json:"connected"` + LastEventAt *float64 `json:"lastEventAt"` + LastInboundAt *float64 `json:"lastInboundAt"` + LastOutboundAt *float64 `json:"lastOutboundAt"` + LastError json.RawMessage `json:"lastError"` + ReconnectAttempts *int `json:"reconnectAttempts"` + Mode *string `json:"mode"` +} + +func (w wireAccount) toAccountState() AccountState { + return AccountState{ + AccountID: strOr(w.AccountID, "default"), + Enabled: boolOr(w.Enabled, true), + Configured: boolOr(w.Configured, true), + Running: boolOr(w.Running, false), + Connected: boolOr(w.Connected, false), + Mode: strOr(w.Mode, ""), + LastEventAt: msToTime(w.LastEventAt), + LastInboundAt: msToTime(w.LastInboundAt), + LastOutboundAt: msToTime(w.LastOutboundAt), + LastError: lastErrorString(w.LastError), + ReconnectAttempts: intOr(w.ReconnectAttempts, 0), + } +} + +func strOr(p *string, def string) string { + if p == nil { + return def + } + return *p +} + +func boolOr(p *bool, def bool) bool { + if p == nil { + return def + } + return *p +} + +func intOr(p *int, def int) int { + if p == nil { + return def + } + return *p +} + +// msToTime converts a Unix-milliseconds timestamp to *time.Time. Zero and +// negative values are treated as absent. +func msToTime(p *float64) *time.Time { + if p == nil || *p <= 0 { + return nil + } + t := time.UnixMilli(int64(*p)).UTC() + return &t +} + +// lastErrorString normalizes the lastError field, which may be absent, +// null, a plain string, or an object, into a display string. +func lastErrorString(raw json.RawMessage) string { + trimmed := bytes.TrimSpace(raw) + if len(trimmed) == 0 || string(trimmed) == "null" { + return "" + } + var s string + if err := json.Unmarshal(trimmed, &s); err == nil { + return s + } + var obj map[string]any + if err := json.Unmarshal(trimmed, &obj); err == nil { + for _, key := range []string{"message", "error", "reason", "code"} { + if v, ok := obj[key].(string); ok && v != "" { + return v + } + } + } + return string(trimmed) +} + +// BuildChannelStates parses a channels.status payload and evaluates the +// health of every channel account. The result is sorted by (channel, +// account_id) for deterministic output. +func BuildChannelStates(payload []byte, now time.Time) ([]ChannelState, error) { + var wp wirePayload + if err := json.Unmarshal(payload, &wp); err != nil { + return nil, fmt.Errorf("parse channels.status payload: %w", err) + } + + states := make([]ChannelState, 0, len(wp.ChannelAccounts)) + for channel, accounts := range wp.ChannelAccounts { + for _, wa := range accounts { + a := wa.toAccountState() + states = append(states, ChannelState{ + Channel: channel, + AccountID: a.AccountID, + Status: EvaluateAccount(a, now), + Enabled: a.Enabled, + Running: a.Running, + Connected: a.Connected, + Mode: a.Mode, + LastEventAt: a.LastEventAt, + LastInboundAt: a.LastInboundAt, + LastOutboundAt: a.LastOutboundAt, + LastError: a.LastError, + ReconnectAttempts: a.ReconnectAttempts, + CheckedAt: now, + }) + } + } + sort.Slice(states, func(i, j int) bool { + if states[i].Channel != states[j].Channel { + return states[i].Channel < states[j].Channel + } + return states[i].AccountID < states[j].AccountID + }) + return states, nil +} + +// isPersistentMode reports whether a channel mode implies a long-lived +// socket connection over which events are expected to keep flowing. +// Pull-style modes (http, webhook) never go "stale". +func isPersistentMode(mode string) bool { + return mode != "http" && mode != "webhook" +} + +// EvaluateAccount derives the health status for one channel account. +func EvaluateAccount(a AccountState, now time.Time) string { + if !a.Enabled || !a.Configured { + return StatusDisabled + } + switch { + case a.Running && !a.Connected: + return StatusDisconnected + case !a.Running: + return StatusNotRunning + case a.Connected: + if isPersistentMode(a.Mode) && a.LastEventAt != nil && now.Sub(*a.LastEventAt) > StaleThreshold { + return StatusStale + } + return StatusHealthy + default: + return StatusUnknown + } +} + +// DeriveOverall computes the instance-level status from the per-channel +// statuses. checked is false when the instance has never been polled. +func DeriveOverall(gatewayReachable, checked bool, channels []ChannelState) string { + if !checked { + return OverallUnknown + } + if !gatewayReachable { + return OverallUnreachable + } + var active, unhealthy, degraded, healthy int + for _, c := range channels { + switch c.Status { + case StatusDisabled: + continue + case StatusDisconnected, StatusNotRunning: + unhealthy++ + case StatusStale, StatusUnknown: + degraded++ + case StatusHealthy: + healthy++ + } + active++ + } + switch { + case active == 0: + return OverallNoChannels + case unhealthy > 0: + return OverallUnhealthy + case degraded > 0: + return OverallDegraded + case healthy > 0: + return OverallHealthy + default: + return OverallUnknown + } +} diff --git a/control-plane/internal/channelhealth/channelhealth_test.go b/control-plane/internal/channelhealth/channelhealth_test.go new file mode 100644 index 00000000..eedbd1dd --- /dev/null +++ b/control-plane/internal/channelhealth/channelhealth_test.go @@ -0,0 +1,283 @@ +package channelhealth + +import ( + "fmt" + "testing" + "time" +) + +var testNow = time.Date(2026, 8, 6, 12, 0, 0, 0, time.UTC) + +func tp(t time.Time) *time.Time { return &t } + +func TestEvaluateAccount(t *testing.T) { + fresh := tp(testNow.Add(-1 * time.Minute)) + old := tp(testNow.Add(-31 * time.Minute)) + + tests := []struct { + name string + acc AccountState + want string + }{ + { + name: "disabled account", + acc: AccountState{Enabled: false, Configured: true, Running: true, Connected: true}, + want: StatusDisabled, + }, + { + name: "unconfigured account", + acc: AccountState{Enabled: true, Configured: false, Running: true, Connected: true}, + want: StatusDisabled, + }, + { + name: "running but not connected", + acc: AccountState{Enabled: true, Configured: true, Running: true, Connected: false}, + want: StatusDisconnected, + }, + { + name: "enabled but not running", + acc: AccountState{Enabled: true, Configured: true, Running: false, Connected: false}, + want: StatusNotRunning, + }, + { + name: "connected socket mode with recent event", + acc: AccountState{Enabled: true, Configured: true, Running: true, Connected: true, Mode: "socket", LastEventAt: fresh}, + want: StatusHealthy, + }, + { + name: "connected socket mode with stale event", + acc: AccountState{Enabled: true, Configured: true, Running: true, Connected: true, Mode: "socket", LastEventAt: old}, + want: StatusStale, + }, + { + name: "connected http mode with old event is not stale", + acc: AccountState{Enabled: true, Configured: true, Running: true, Connected: true, Mode: "http", LastEventAt: old}, + want: StatusHealthy, + }, + { + name: "connected webhook mode with old event is not stale", + acc: AccountState{Enabled: true, Configured: true, Running: true, Connected: true, Mode: "webhook", LastEventAt: old}, + want: StatusHealthy, + }, + { + name: "connected persistent mode without lastEventAt", + acc: AccountState{Enabled: true, Configured: true, Running: true, Connected: true, Mode: "socket"}, + want: StatusHealthy, + }, + { + name: "stale boundary: exactly at threshold is not stale", + acc: AccountState{Enabled: true, Configured: true, Running: true, Connected: true, Mode: "socket", LastEventAt: tp(testNow.Add(-StaleThreshold))}, + want: StatusHealthy, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := EvaluateAccount(tt.acc, testNow); got != tt.want { + t.Errorf("EvaluateAccount() = %q, want %q", got, tt.want) + } + }) + } +} + +func TestBuildChannelStates(t *testing.T) { + now := time.Date(2026, 8, 6, 12, 1, 0, 0, time.UTC) + wantEvent := time.Date(2026, 8, 6, 11, 59, 0, 0, time.UTC) + eventMs := wantEvent.UnixMilli() + + payload := []byte(fmt.Sprintf(`{ + "ts": %d, + "channelOrder": ["slack", "telegram"], + "channels": {"slack": {"configured": true}, "telegram": {"configured": true}}, + "channelAccounts": { + "slack": [{ + "accountId": "default", + "enabled": true, + "configured": true, + "running": true, + "connected": true, + "lastConnectedAt": %d, + "lastEventAt": %d, + "lastInboundAt": %d, + "lastError": null, + "reconnectAttempts": 0, + "mode": "socket", + "restartPending": false + }], + "telegram": [{ + "accountId": "bot1", + "enabled": true, + "configured": true, + "running": true, + "connected": false, + "lastError": {"message": "invalid token", "code": "AUTH"}, + "reconnectAttempts": 3, + "mode": "socket" + }, { + "accountId": "bot2", + "enabled": false, + "configured": true, + "running": false, + "connected": false + }] + }, + "channelDefaultAccountId": {"slack": "default"} + }`, now.UnixMilli(), wantEvent.Add(-1*time.Hour).UnixMilli(), eventMs, eventMs)) + + states, err := BuildChannelStates(payload, now) + if err != nil { + t.Fatalf("BuildChannelStates() error: %v", err) + } + if len(states) != 3 { + t.Fatalf("got %d states, want 3", len(states)) + } + + // Sorted by (channel, account_id): slack/default, telegram/bot1, telegram/bot2. + slack := states[0] + if slack.Channel != "slack" || slack.AccountID != "default" { + t.Fatalf("states[0] = %s/%s, want slack/default", slack.Channel, slack.AccountID) + } + if slack.Status != StatusHealthy { + t.Errorf("slack status = %q, want %q", slack.Status, StatusHealthy) + } + if slack.LastEventAt == nil || !slack.LastEventAt.Equal(wantEvent) { + t.Errorf("slack.LastEventAt = %v, want %v", slack.LastEventAt, wantEvent) + } + if slack.LastInboundAt == nil || !slack.LastInboundAt.Equal(wantEvent) { + t.Errorf("slack.LastInboundAt = %v, want %v", slack.LastInboundAt, wantEvent) + } + if slack.LastOutboundAt != nil { + t.Errorf("slack.LastOutboundAt = %v, want nil (absent field)", slack.LastOutboundAt) + } + if slack.LastError != "" { + t.Errorf("slack.LastError = %q, want empty (null on wire)", slack.LastError) + } + if !slack.CheckedAt.Equal(now) { + t.Errorf("slack.CheckedAt = %v, want %v", slack.CheckedAt, now) + } + + bot1 := states[1] + if bot1.Channel != "telegram" || bot1.AccountID != "bot1" { + t.Fatalf("states[1] = %s/%s, want telegram/bot1", bot1.Channel, bot1.AccountID) + } + if bot1.Status != StatusDisconnected { + t.Errorf("bot1 status = %q, want %q", bot1.Status, StatusDisconnected) + } + if bot1.LastError != "invalid token" { + t.Errorf("bot1.LastError = %q, want %q (object message extracted)", bot1.LastError, "invalid token") + } + if bot1.ReconnectAttempts != 3 { + t.Errorf("bot1.ReconnectAttempts = %d, want 3", bot1.ReconnectAttempts) + } + + bot2 := states[2] + if bot2.Status != StatusDisabled { + t.Errorf("bot2 status = %q, want %q (disabled accounts are recorded)", bot2.Status, StatusDisabled) + } +} + +func TestBuildChannelStatesDefensive(t *testing.T) { + t.Run("missing fields default sensibly", func(t *testing.T) { + payload := []byte(`{"channelAccounts": {"slack": [{}]}}`) + states, err := BuildChannelStates(payload, testNow) + if err != nil { + t.Fatalf("error: %v", err) + } + if len(states) != 1 { + t.Fatalf("got %d states, want 1", len(states)) + } + s := states[0] + if s.AccountID != "default" { + t.Errorf("AccountID = %q, want %q", s.AccountID, "default") + } + // enabled/configured default true, running/connected default false + // => enabled && !running => not_running. + if s.Status != StatusNotRunning { + t.Errorf("Status = %q, want %q", s.Status, StatusNotRunning) + } + if s.LastEventAt != nil || s.LastInboundAt != nil || s.LastOutboundAt != nil { + t.Errorf("timestamps should be nil for absent fields") + } + }) + + t.Run("lastError as plain string", func(t *testing.T) { + payload := []byte(`{"channelAccounts": {"slack": [{"running": true, "lastError": "boom"}]}}`) + states, err := BuildChannelStates(payload, testNow) + if err != nil { + t.Fatalf("error: %v", err) + } + if states[0].LastError != "boom" { + t.Errorf("LastError = %q, want %q", states[0].LastError, "boom") + } + }) + + t.Run("lastError object without message falls back to raw JSON", func(t *testing.T) { + payload := []byte(`{"channelAccounts": {"slack": [{"lastError": {"weird": 1}}]}}`) + states, err := BuildChannelStates(payload, testNow) + if err != nil { + t.Fatalf("error: %v", err) + } + if states[0].LastError != `{"weird": 1}` { + t.Errorf("LastError = %q, want raw JSON", states[0].LastError) + } + }) + + t.Run("empty payload", func(t *testing.T) { + states, err := BuildChannelStates([]byte(`{}`), testNow) + if err != nil { + t.Fatalf("error: %v", err) + } + if len(states) != 0 { + t.Errorf("got %d states, want 0", len(states)) + } + }) + + t.Run("invalid payload", func(t *testing.T) { + if _, err := BuildChannelStates([]byte(`not json`), testNow); err == nil { + t.Error("expected error for invalid JSON") + } + }) + + t.Run("zero-ms timestamp treated as absent", func(t *testing.T) { + payload := []byte(`{"channelAccounts": {"slack": [{"lastEventAt": 0}]}}`) + states, err := BuildChannelStates(payload, testNow) + if err != nil { + t.Fatalf("error: %v", err) + } + if states[0].LastEventAt != nil { + t.Errorf("LastEventAt = %v, want nil for 0", states[0].LastEventAt) + } + }) +} + +func TestDeriveOverall(t *testing.T) { + ch := func(status string) ChannelState { return ChannelState{Status: status} } + + tests := []struct { + name string + reachable bool + checked bool + channels []ChannelState + want string + }{ + {"never checked", false, false, nil, OverallUnknown}, + {"gateway unreachable", false, true, []ChannelState{ch(StatusHealthy)}, OverallUnreachable}, + {"no channels", true, true, nil, OverallNoChannels}, + {"only disabled channels", true, true, []ChannelState{ch(StatusDisabled)}, OverallNoChannels}, + {"one disconnected", true, true, []ChannelState{ch(StatusHealthy), ch(StatusDisconnected)}, OverallUnhealthy}, + {"one not_running", true, true, []ChannelState{ch(StatusHealthy), ch(StatusNotRunning)}, OverallUnhealthy}, + {"disconnected trumps stale", true, true, []ChannelState{ch(StatusStale), ch(StatusDisconnected)}, OverallUnhealthy}, + {"one stale", true, true, []ChannelState{ch(StatusHealthy), ch(StatusStale)}, OverallDegraded}, + {"one unknown", true, true, []ChannelState{ch(StatusHealthy), ch(StatusUnknown)}, OverallDegraded}, + {"all healthy", true, true, []ChannelState{ch(StatusHealthy), ch(StatusHealthy)}, OverallHealthy}, + {"healthy plus disabled", true, true, []ChannelState{ch(StatusHealthy), ch(StatusDisabled)}, OverallHealthy}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := DeriveOverall(tt.reachable, tt.checked, tt.channels); got != tt.want { + t.Errorf("DeriveOverall() = %q, want %q", got, tt.want) + } + }) + } +} diff --git a/control-plane/internal/channelhealth/monitor.go b/control-plane/internal/channelhealth/monitor.go new file mode 100644 index 00000000..39c40c2b --- /dev/null +++ b/control-plane/internal/channelhealth/monitor.go @@ -0,0 +1,363 @@ +package channelhealth + +import ( + "context" + "encoding/json" + "fmt" + "log" + "sync" + "time" + + "github.com/coder/websocket" + "github.com/gluk-w/claworc/control-plane/internal/database" + "github.com/gluk-w/claworc/control-plane/internal/sshproxy" + "github.com/gluk-w/claworc/control-plane/internal/utils" + "gorm.io/gorm/clause" +) + +const ( + // maxConcurrentChecks bounds how many instances are polled in parallel. + maxConcurrentChecks = 5 + // perInstanceTimeout bounds one instance's dial+RPC round trip. + perInstanceTimeout = 15 * time.Second + // gatewayTunnelLabel is the tunnel manager's label for the OpenClaw + // gateway tunnel (see sshproxy tunnel provisioning). + gatewayTunnelLabel = "Gateway" +) + +// Monitor periodically polls the OpenClaw gateway of every running +// instance for channel health, persists the results, and keeps an +// in-memory snapshot per instance for cheap reads by handlers. +type Monitor struct { + tunnels *sshproxy.TunnelManager + interval time.Duration + + mu sync.RWMutex + snapshots map[uint]Snapshot +} + +// New builds a Monitor. tunnels is used to resolve the local port of each +// instance's Gateway SSH tunnel; interval is the polling period (<=0 falls +// back to 60s). +func New(tunnels *sshproxy.TunnelManager, interval time.Duration) *Monitor { + if interval <= 0 { + interval = 60 * time.Second + } + return &Monitor{ + tunnels: tunnels, + interval: interval, + snapshots: make(map[uint]Snapshot), + } +} + +// Start launches the background polling loop. It returns immediately; the +// goroutine exits when ctx is canceled. +func (m *Monitor) Start(ctx context.Context) { + go m.loop(ctx) +} + +func (m *Monitor) loop(ctx context.Context) { + t := time.NewTicker(m.interval) + defer t.Stop() + + // Run once on startup so the UI has data quickly. + m.checkAll(ctx) + + for { + select { + case <-ctx.Done(): + return + case <-t.C: + m.checkAll(ctx) + } + } +} + +func (m *Monitor) checkAll(ctx context.Context) { + var instances []database.Instance + if err := database.DB.Where("status = ?", "running").Find(&instances).Error; err != nil { + log.Printf("[channelhealth] list instances: %v", err) + return + } + + sem := make(chan struct{}, maxConcurrentChecks) + var wg sync.WaitGroup + for i := range instances { + if ctx.Err() != nil { + break + } + inst := instances[i] + wg.Add(1) + sem <- struct{}{} + go func() { + defer wg.Done() + defer func() { <-sem }() + cctx, cancel := context.WithTimeout(ctx, perInstanceTimeout) + defer cancel() + m.checkInstance(cctx, inst) + }() + } + wg.Wait() +} + +func (m *Monitor) checkInstance(ctx context.Context, inst database.Instance) { + now := time.Now().UTC() + + port, ok := m.gatewayPort(inst.ID) + if !ok { + m.recordUnreachable(inst.ID, now) + return + } + + var gatewayToken string + if inst.GatewayToken != "" { + if tok, err := utils.Decrypt(inst.GatewayToken); err == nil { + gatewayToken = tok + } + } + + payload, err := queryChannelsStatus(ctx, port, gatewayToken) + if err != nil { + log.Printf("[channelhealth] instance %d: channels.status: %v", inst.ID, err) + m.recordUnreachable(inst.ID, now) + return + } + + states, err := BuildChannelStates(payload, now) + if err != nil { + // The gateway responded but with an unparseable payload; keep the + // previous snapshot rather than flapping to unreachable. + log.Printf("[channelhealth] instance %d: %v", inst.ID, err) + return + } + + if err := persistStates(inst.ID, states); err != nil { + log.Printf("[channelhealth] instance %d: persist: %v", inst.ID, err) + } + + m.store(Snapshot{ + InstanceID: inst.ID, + Overall: DeriveOverall(true, true, states), + GatewayReachable: true, + CheckedAt: now, + Channels: states, + }) +} + +// gatewayPort resolves the local port of the instance's active Gateway +// tunnel. +func (m *Monitor) gatewayPort(instanceID uint) (int, bool) { + if m.tunnels == nil { + return 0, false + } + for _, t := range m.tunnels.GetTunnelsForInstance(instanceID) { + if t.Label == gatewayTunnelLabel && t.Status == "active" { + return t.LocalPort, true + } + } + return 0, false +} + +// recordUnreachable marks the instance-level state unreachable while +// keeping the previously known channel rows (from the prior snapshot or, +// failing that, the database) so the UI can still show the last state. +func (m *Monitor) recordUnreachable(instanceID uint, now time.Time) { + m.mu.RLock() + prev, had := m.snapshots[instanceID] + m.mu.RUnlock() + + channels := prev.Channels + if !had { + if dbSnap, ok := SnapshotFromDB(instanceID); ok { + channels = dbSnap.Channels + } + } + + m.store(Snapshot{ + InstanceID: instanceID, + Overall: OverallUnreachable, + GatewayReachable: false, + CheckedAt: now, + Channels: channels, + }) +} + +// store swaps in the new snapshot and logs notable overall-status +// transitions (to unhealthy/unreachable, and recovery back to healthy). +func (m *Monitor) store(snap Snapshot) { + m.mu.Lock() + prev, had := m.snapshots[snap.InstanceID] + m.snapshots[snap.InstanceID] = snap + m.mu.Unlock() + + prevOverall := OverallUnknown + if had { + prevOverall = prev.Overall + } + if prevOverall == snap.Overall { + return + } + switch { + case snap.Overall == OverallUnhealthy || snap.Overall == OverallUnreachable: + log.Printf("[channelhealth] instance %d: channel health %s -> %s", snap.InstanceID, prevOverall, snap.Overall) + case snap.Overall == OverallHealthy && (prevOverall == OverallUnhealthy || prevOverall == OverallUnreachable): + log.Printf("[channelhealth] instance %d: channel health recovered: %s -> %s", snap.InstanceID, prevOverall, snap.Overall) + } +} + +// Snapshot returns a copy of the latest snapshot for the instance, or +// ok=false when the instance has never been checked. +func (m *Monitor) Snapshot(instanceID uint) (Snapshot, bool) { + m.mu.RLock() + snap, ok := m.snapshots[instanceID] + m.mu.RUnlock() + if !ok { + return Snapshot{}, false + } + out := snap + out.Channels = append([]ChannelState(nil), snap.Channels...) + return out, true +} + +// SnapshotFromDB reconstructs a snapshot from persisted rows. Used as a +// fallback when no in-memory snapshot exists yet (e.g. right after a +// control-plane restart). ok=false when no rows exist. +func SnapshotFromDB(instanceID uint) (Snapshot, bool) { + var rows []database.ChannelHealthStatus + if err := database.DB.Where("instance_id = ?", instanceID). + Order("channel ASC, account_id ASC").Find(&rows).Error; err != nil || len(rows) == 0 { + return Snapshot{}, false + } + + channels := make([]ChannelState, len(rows)) + var latest time.Time + for i, r := range rows { + channels[i] = ChannelState{ + Channel: r.Channel, + AccountID: r.AccountID, + Status: r.Status, + Enabled: r.Enabled, + Running: r.Running, + Connected: r.Connected, + Mode: r.Mode, + LastEventAt: r.LastEventAt, + LastInboundAt: r.LastInboundAt, + LastOutboundAt: r.LastOutboundAt, + LastError: r.LastError, + ReconnectAttempts: r.ReconnectAttempts, + CheckedAt: r.CheckedAt, + } + if r.CheckedAt.After(latest) { + latest = r.CheckedAt + } + } + + return Snapshot{ + InstanceID: instanceID, + Overall: DeriveOverall(true, true, channels), + GatewayReachable: true, + CheckedAt: latest, + Channels: channels, + }, true +} + +// persistStates upserts one row per (instance, channel, account) and +// deletes rows for accounts that disappeared from the gateway's config. +func persistStates(instanceID uint, states []ChannelState) error { + var existing []database.ChannelHealthStatus + if err := database.DB.Where("instance_id = ?", instanceID).Find(&existing).Error; err != nil { + return err + } + keep := make(map[[2]string]bool, len(states)) + for _, s := range states { + keep[[2]string{s.Channel, s.AccountID}] = true + } + for _, e := range existing { + if !keep[[2]string{e.Channel, e.AccountID}] { + if err := database.DB.Delete(&database.ChannelHealthStatus{}, e.ID).Error; err != nil { + return err + } + } + } + + for _, s := range states { + row := database.ChannelHealthStatus{ + InstanceID: instanceID, + Channel: s.Channel, + AccountID: s.AccountID, + Status: s.Status, + Enabled: s.Enabled, + Running: s.Running, + Connected: s.Connected, + Mode: s.Mode, + LastEventAt: s.LastEventAt, + LastInboundAt: s.LastInboundAt, + LastOutboundAt: s.LastOutboundAt, + LastError: s.LastError, + ReconnectAttempts: s.ReconnectAttempts, + CheckedAt: s.CheckedAt, + } + err := database.DB.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "instance_id"}, {Name: "channel"}, {Name: "account_id"}}, + DoUpdates: clause.AssignmentColumns([]string{ + "status", "enabled", "running", "connected", "mode", + "last_event_at", "last_inbound_at", "last_outbound_at", + "last_error", "reconnect_attempts", "checked_at", "updated_at", + }), + }).Create(&row).Error + if err != nil { + return err + } + } + return nil +} + +// queryChannelsStatus dials the gateway over the local tunnel port, issues +// a channels.status request (without probing — probes hit provider APIs), +// and returns the raw response payload. +func queryChannelsStatus(ctx context.Context, port int, gatewayToken string) (json.RawMessage, error) { + conn, err := sshproxy.DialGateway(ctx, port, gatewayToken) + if err != nil { + return nil, err + } + defer conn.CloseNow() + + reqID := fmt.Sprintf("chanhealth-%d", time.Now().UnixNano()) + frame := map[string]any{ + "type": "req", + "id": reqID, + "method": "channels.status", + "params": map[string]any{}, + } + reqJSON, err := json.Marshal(frame) + if err != nil { + return nil, fmt.Errorf("marshal channels.status: %w", err) + } + if err := conn.Write(ctx, websocket.MessageText, reqJSON); err != nil { + return nil, fmt.Errorf("send channels.status: %w", err) + } + + for { + _, data, err := conn.Read(ctx) + if err != nil { + return nil, fmt.Errorf("read channels.status: %w", err) + } + var resp struct { + Type string `json:"type"` + ID string `json:"id"` + OK bool `json:"ok"` + Payload json.RawMessage `json:"payload"` + Error json.RawMessage `json:"error"` + } + if err := json.Unmarshal(data, &resp); err != nil { + continue + } + if resp.Type != "res" || resp.ID != reqID { + continue + } + if !resp.OK { + return nil, fmt.Errorf("channels.status failed: %s", string(resp.Error)) + } + return resp.Payload, nil + } +} diff --git a/control-plane/internal/config/config.go b/control-plane/internal/config/config.go index be268129..0b633245 100644 --- a/control-plane/internal/config/config.go +++ b/control-plane/internal/config/config.go @@ -47,6 +47,13 @@ type Settings struct { // on every frame received, so an actively-streaming agent is never cut off; // only a genuine stall trips it. WebhookIdleTimeout time.Duration `envconfig:"WEBHOOK_IDLE_TIMEOUT" default:"120s"` + + // Channel health monitor settings. The monitor polls each running + // instance's OpenClaw gateway (channels.status) and exposes per-channel + // health via the API. Disabling it removes the background polling and + // makes the channel health endpoints report "disabled". + ChannelHealthEnabled bool `envconfig:"CHANNEL_HEALTH_ENABLED" default:"true"` + ChannelHealthInterval time.Duration `envconfig:"CHANNEL_HEALTH_INTERVAL" default:"60s"` } var Cfg Settings diff --git a/control-plane/internal/database/migrations/migration_00001_baseline.go b/control-plane/internal/database/migrations/migration_00001_baseline.go index 939118cc..4c9a0db0 100644 --- a/control-plane/internal/database/migrations/migration_00001_baseline.go +++ b/control-plane/internal/database/migrations/migration_00001_baseline.go @@ -63,5 +63,6 @@ func AutoMigrateAll(gdb interface { &models.TeamProvider{}, &models.WebhookApiKey{}, &models.WebhookLog{}, + &models.ChannelHealthStatus{}, ) } diff --git a/control-plane/internal/database/models.go b/control-plane/internal/database/models.go index ac84b8f8..ed475d62 100644 --- a/control-plane/internal/database/models.go +++ b/control-plane/internal/database/models.go @@ -9,32 +9,33 @@ import "github.com/gluk-w/claworc/control-plane/internal/database/models" // types via the GORM Migrator without an import cycle. type ( - Skill = models.Skill - Instance = models.Instance - Team = models.Team - TeamMember = models.TeamMember - TeamProvider = models.TeamProvider - BrowserSession = models.BrowserSession - ProviderModel = models.ProviderModel - ProviderModelCost = models.ProviderModelCost - LLMProvider = models.LLMProvider - LLMGatewayKey = models.LLMGatewayKey - LLMRequestLog = models.LLMRequestLog - Setting = models.Setting - User = models.User - UserInstance = models.UserInstance - Backup = models.Backup - BackupSchedule = models.BackupSchedule - SharedFolder = models.SharedFolder - KanbanBoard = models.KanbanBoard - KanbanTask = models.KanbanTask - KanbanComment = models.KanbanComment - KanbanArtifact = models.KanbanArtifact - InstanceSoul = models.InstanceSoul - WebAuthnCredential = models.WebAuthnCredential - UserSSHKey = models.UserSSHKey - WebhookApiKey = models.WebhookApiKey - WebhookLog = models.WebhookLog + Skill = models.Skill + Instance = models.Instance + Team = models.Team + TeamMember = models.TeamMember + TeamProvider = models.TeamProvider + BrowserSession = models.BrowserSession + ProviderModel = models.ProviderModel + ProviderModelCost = models.ProviderModelCost + LLMProvider = models.LLMProvider + LLMGatewayKey = models.LLMGatewayKey + LLMRequestLog = models.LLMRequestLog + Setting = models.Setting + User = models.User + UserInstance = models.UserInstance + Backup = models.Backup + BackupSchedule = models.BackupSchedule + SharedFolder = models.SharedFolder + KanbanBoard = models.KanbanBoard + KanbanTask = models.KanbanTask + KanbanComment = models.KanbanComment + KanbanArtifact = models.KanbanArtifact + InstanceSoul = models.InstanceSoul + WebAuthnCredential = models.WebAuthnCredential + UserSSHKey = models.UserSSHKey + WebhookApiKey = models.WebhookApiKey + WebhookLog = models.WebhookLog + ChannelHealthStatus = models.ChannelHealthStatus ) // Helper re-exports keep `database.ParseTeamIDs(...)` etc. working for diff --git a/control-plane/internal/database/models/models.go b/control-plane/internal/database/models/models.go index f808ca8d..5c35d57f 100644 --- a/control-plane/internal/database/models/models.go +++ b/control-plane/internal/database/models/models.go @@ -448,6 +448,29 @@ type WebAuthnCredential struct { CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` } +// ChannelHealthStatus is the latest observed runtime health for one +// channel account on one instance, refreshed periodically by the +// channelhealth monitor from the gateway's channels.status RPC. +type ChannelHealthStatus struct { + ID uint `gorm:"primaryKey;autoIncrement" json:"id"` + InstanceID uint `gorm:"not null;index:idx_channel_health_instance_channel_account,unique" json:"instance_id"` + Channel string `gorm:"not null;index:idx_channel_health_instance_channel_account,unique" json:"channel"` + AccountID string `gorm:"not null;default:'';index:idx_channel_health_instance_channel_account,unique" json:"account_id"` + Status string `gorm:"not null;default:''" json:"status"` // healthy|disconnected|not_running|stale|disabled|unknown + Enabled bool `gorm:"not null;default:false" json:"enabled"` + Running bool `gorm:"not null;default:false" json:"running"` + Connected bool `gorm:"not null;default:false" json:"connected"` + Mode string `gorm:"default:''" json:"mode"` + LastEventAt *time.Time `json:"last_event_at"` + LastInboundAt *time.Time `json:"last_inbound_at"` + LastOutboundAt *time.Time `json:"last_outbound_at"` + LastError string `gorm:"type:text;default:''" json:"last_error"` + ReconnectAttempts int `gorm:"default:0" json:"reconnect_attempts"` + CheckedAt time.Time `json:"checked_at"` + CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` + UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"` +} + // UserSSHKey is a public key a user authenticates with against the inbound // SSH gateway. The private key is never stored — it is generated on demand // and handed to the user exactly once (or the user uploads their own pubkey). diff --git a/control-plane/internal/handlers/channels.go b/control-plane/internal/handlers/channels.go new file mode 100644 index 00000000..fde9e7e0 --- /dev/null +++ b/control-plane/internal/handlers/channels.go @@ -0,0 +1,127 @@ +package handlers + +import ( + "net/http" + "strconv" + "time" + + "github.com/gluk-w/claworc/control-plane/internal/channelhealth" + "github.com/gluk-w/claworc/control-plane/internal/database" + "github.com/gluk-w/claworc/control-plane/internal/middleware" + "github.com/go-chi/chi/v5" +) + +// ChannelHealthMon is set from main.go during init when the channel health +// monitor is enabled. nil means CLAWORC_CHANNEL_HEALTH_ENABLED=false. +var ChannelHealthMon *channelhealth.Monitor + +type channelHealthResponse struct { + InstanceID uint `json:"instance_id"` + Overall string `json:"overall"` + GatewayReachable bool `json:"gateway_reachable"` + CheckedAt *string `json:"checked_at"` + Channels []channelHealthEntry `json:"channels"` +} + +type channelHealthEntry struct { + Channel string `json:"channel"` + AccountID string `json:"account_id"` + Status string `json:"status"` + Enabled bool `json:"enabled"` + Running bool `json:"running"` + Connected bool `json:"connected"` + Mode string `json:"mode"` + LastEventAt *string `json:"last_event_at"` + LastInboundAt *string `json:"last_inbound_at"` + LastOutboundAt *string `json:"last_outbound_at"` + LastError string `json:"last_error"` + ReconnectAttempts int `json:"reconnect_attempts"` + CheckedAt string `json:"checked_at"` +} + +// GetChannelHealth returns the latest per-channel health for an instance, +// as observed by the background channel health monitor. +func GetChannelHealth(w http.ResponseWriter, r *http.Request) { + id, err := strconv.Atoi(chi.URLParam(r, "id")) + if err != nil { + writeError(w, http.StatusBadRequest, "Invalid instance ID") + return + } + + var inst database.Instance + if err := database.DB.First(&inst, id).Error; err != nil { + writeError(w, http.StatusNotFound, "Instance not found") + return + } + + if !middleware.CanAccessInstance(r, inst.ID) { + writeError(w, http.StatusForbidden, "Access denied") + return + } + + // Monitor disabled entirely (CLAWORC_CHANNEL_HEALTH_ENABLED=false). + if ChannelHealthMon == nil { + writeJSON(w, http.StatusOK, channelHealthResponse{ + InstanceID: inst.ID, + Overall: "disabled", + GatewayReachable: false, + CheckedAt: nil, + Channels: []channelHealthEntry{}, + }) + return + } + + snap, ok := ChannelHealthMon.Snapshot(inst.ID) + if !ok { + // No in-memory snapshot yet (e.g. control plane just restarted): + // fall back to persisted rows. + snap, ok = channelhealth.SnapshotFromDB(inst.ID) + } + if !ok { + writeJSON(w, http.StatusOK, channelHealthResponse{ + InstanceID: inst.ID, + Overall: channelhealth.OverallUnknown, + GatewayReachable: false, + CheckedAt: nil, + Channels: []channelHealthEntry{}, + }) + return + } + + channels := make([]channelHealthEntry, len(snap.Channels)) + for i, c := range snap.Channels { + channels[i] = channelHealthEntry{ + Channel: c.Channel, + AccountID: c.AccountID, + Status: c.Status, + Enabled: c.Enabled, + Running: c.Running, + Connected: c.Connected, + Mode: c.Mode, + LastEventAt: rfc3339OrNil(c.LastEventAt), + LastInboundAt: rfc3339OrNil(c.LastInboundAt), + LastOutboundAt: rfc3339OrNil(c.LastOutboundAt), + LastError: c.LastError, + ReconnectAttempts: c.ReconnectAttempts, + CheckedAt: c.CheckedAt.UTC().Format(time.RFC3339), + } + } + + checkedAt := snap.CheckedAt.UTC().Format(time.RFC3339) + writeJSON(w, http.StatusOK, channelHealthResponse{ + InstanceID: inst.ID, + Overall: snap.Overall, + GatewayReachable: snap.GatewayReachable, + CheckedAt: &checkedAt, + Channels: channels, + }) +} + +// rfc3339OrNil formats an optional timestamp as RFC3339 UTC, or nil. +func rfc3339OrNil(t *time.Time) *string { + if t == nil || t.IsZero() { + return nil + } + s := t.UTC().Format(time.RFC3339) + return &s +} diff --git a/control-plane/internal/handlers/instances.go b/control-plane/internal/handlers/instances.go index dfa14e25..e97808a0 100644 --- a/control-plane/internal/handlers/instances.go +++ b/control-plane/internal/handlers/instances.go @@ -15,6 +15,7 @@ import ( "time" "github.com/gluk-w/claworc/control-plane/internal/analytics" + "github.com/gluk-w/claworc/control-plane/internal/channelhealth" "github.com/gluk-w/claworc/control-plane/internal/config" "github.com/gluk-w/claworc/control-plane/internal/database" "github.com/gluk-w/claworc/control-plane/internal/llmgateway" @@ -191,6 +192,40 @@ type instanceResponse struct { Affinity string `json:"affinity"` ServiceAccountAnnotations map[string]string `json:"service_account_annotations"` Ports []orchestrator.PortSpec `json:"ports"` + ChannelHealth *channelHealthSummary `json:"channel_health"` +} + +// channelHealthSummary is the compact channel-health view embedded in +// instance responses. nil (JSON null) when the monitor is disabled or has +// no data for the instance yet. +type channelHealthSummary struct { + Overall string `json:"overall"` + UnhealthyCount int `json:"unhealthy_count"` + CheckedAt string `json:"checked_at"` +} + +// channelHealthSummaryFor reads the monitor's in-memory snapshot only — +// never the database — so it stays cheap on the instance list path. +func channelHealthSummaryFor(instanceID uint) *channelHealthSummary { + if ChannelHealthMon == nil { + return nil + } + snap, ok := ChannelHealthMon.Snapshot(instanceID) + if !ok { + return nil + } + unhealthy := 0 + for _, c := range snap.Channels { + switch c.Status { + case channelhealth.StatusDisconnected, channelhealth.StatusNotRunning, channelhealth.StatusStale: + unhealthy++ + } + } + return &channelHealthSummary{ + Overall: snap.Overall, + UnhealthyCount: unhealthy, + CheckedAt: snap.CheckedAt.UTC().Format(time.RFC3339), + } } func generateName(displayName string) string { @@ -558,6 +593,7 @@ func instanceToResponse(inst database.Instance, status string) instanceResponse Affinity: inst.Affinity, ServiceAccountAnnotations: serviceAccountAnnotations, Ports: ports, + ChannelHealth: channelHealthSummaryFor(inst.ID), } } diff --git a/control-plane/main.go b/control-plane/main.go index fe8238c2..3d5d60ea 100644 --- a/control-plane/main.go +++ b/control-plane/main.go @@ -19,6 +19,7 @@ import ( "github.com/gluk-w/claworc/control-plane/internal/auth" "github.com/gluk-w/claworc/control-plane/internal/backup" "github.com/gluk-w/claworc/control-plane/internal/browserprov" + "github.com/gluk-w/claworc/control-plane/internal/channelhealth" "github.com/gluk-w/claworc/control-plane/internal/config" "github.com/gluk-w/claworc/control-plane/internal/database" "github.com/gluk-w/claworc/control-plane/internal/handlers" @@ -295,6 +296,16 @@ func main() { handlers.ModeratorSvc.StartSummarizer(ctx) } + // Start background channel health monitor. It polls each running + // instance's OpenClaw gateway (channels.status) over the Gateway SSH + // tunnel and exposes results via /instances/{id}/channels/health. + if config.Cfg.ChannelHealthEnabled { + chMon := channelhealth.New(tunnelMgr, config.Cfg.ChannelHealthInterval) + handlers.ChannelHealthMon = chMon + chMon.Start(ctx) + log.Printf("Channel health monitor started (interval=%s)", config.Cfg.ChannelHealthInterval) + } + // Start background SSH key rotation job (checks daily) cancelRotation := handlers.StartKeyRotationJob(ctx) _ = cancelRotation // stopped via context cancellation on shutdown @@ -367,6 +378,7 @@ func main() { r.Get("/instances/{id}/logs", handlers.StreamLogs) r.Get("/instances/{id}/ssh-test", handlers.SSHConnectionTest) r.Get("/instances/{id}/ssh-status", handlers.GetSSHStatus) + r.Get("/instances/{id}/channels/health", handlers.GetChannelHealth) r.Get("/instances/{id}/ssh-events", handlers.GetSSHEvents) r.Post("/instances/{id}/ssh-reconnect", handlers.SSHReconnect) r.Get("/instances/{id}/tunnels", handlers.GetTunnelStatus) diff --git a/docs/README.md b/docs/README.md index 1758588f..de60d130 100644 --- a/docs/README.md +++ b/docs/README.md @@ -27,5 +27,6 @@ Claworc replaces this manual approach by: | [UI](ui.md) | Frontend pages, components, and interaction patterns | | [Environment Variables](environment-variables.md) | Global and per-instance env vars, reserved names, and skill `required_env_vars` | | [SSH Connectivity](ssh-connectivity.md) | SSH architecture, tunnels, health monitoring, and key rotation | +| [Channel Health Monitoring](channel-monitoring.md) | Per-channel liveness monitoring of each instance's chat channels (Slack, Telegram, Discord, …) | | [Kubernetes Deployment](deployment/kubernetes.md) | Kubernetes deployment guide with SSH network policies and security contexts | | [Docker Deployment](deployment/docker.md) | Docker deployment guide with SSH network configuration | diff --git a/docs/channel-monitoring.md b/docs/channel-monitoring.md new file mode 100644 index 00000000..217c477f --- /dev/null +++ b/docs/channel-monitoring.md @@ -0,0 +1,145 @@ +# Channel Health Monitoring + +## Overview + +The control plane monitors whether each OpenClaw instance's chat channels +(Slack, Telegram, Discord, WhatsApp, etc.) are alive and receiving events. +Every 60 seconds (configurable) it connects to the instance's OpenClaw +gateway over the existing SSH tunnel, calls the gateway's `channels.status` +API, evaluates per-channel health, and persists the latest result. The UI +surfaces this as a **Channel Health** panel on the Agent detail page and a +warning indicator in the agent list. + +Monitoring only observes. It never restarts channels or delivers alerts — +see [Future work](#future-work). + +## Architecture + +``` +poller (every CLAWORC_CHANNEL_HEALTH_INTERVAL, default 60s) + → OpenClaw gateway over the existing SSH tunnel (WS RPC `channels.status`) + → health evaluation (per-channel + overall) + → persistence (`channel_health_statuses` table) + in-memory snapshot + → API (`/api/v1/instances/{id}/channels/health`, instance-list summary) + → UI (Channel Health panel, agent-list warning indicator) +``` + +1. A background poller ticks once per interval and checks every running + instance. Checks ride the control plane's existing multiplexed SSH + connection to the instance — no new connections are dialed. +2. Each check calls the OpenClaw gateway's `channels.status` WebSocket RPC, + which reports every channel known to the OpenClaw config along with its + connection state, last-event time, last error, and reconnect count. +3. The evaluator maps the raw gateway report onto the health states below + and computes the instance-level overall status. +4. The latest per-channel status is written to the + `channel_health_statuses` table and kept in an in-memory snapshot; the + API serves from the snapshot and falls back to the persisted rows after + a control-plane restart. + +## Health states + +Per channel: + +| State | Meaning | +|---|---| +| `healthy` | Channel connected and receiving events | +| `stale` | Socket connected but no events for over **30 minutes**. Applies only to persistent-socket modes (e.g. Slack Socket Mode); webhook/http modes are exempt | +| `disconnected` | Channel running but its connection to the provider is down | +| `not_running` | Channel enabled in the OpenClaw config but not running | +| `disabled` | Channel disabled/unconfigured in the OpenClaw config | + +Instance level: + +| State | Meaning | +|---|---| +| `unreachable` | Gateway not responding — the OpenClaw process may be down | +| `no_channels` | Gateway reachable but no channels configured | +| `unknown` | Not yet checked | + +Overall instance status is derived from the per-channel states: +**unhealthy** if any channel is `disconnected` or `not_running`, +**degraded** if any channel is `stale`, **healthy** otherwise. + +## Configuration + +| Env var | Default | Meaning | +|---|---|---| +| `CLAWORC_CHANNEL_HEALTH_ENABLED` | `true` | Enable the poller | +| `CLAWORC_CHANNEL_HEALTH_INTERVAL` | `60s` | Time between checks | + +## API + +`GET /api/v1/instances/{id}/channels/health` returns the overall status +plus per-channel detail: + +```json +{ + "overall": "degraded", + "checked_at": "2026-08-06T10:15:00Z", + "channels": [ + { + "channel": "slack", + "status": "healthy", + "last_event_at": "2026-08-06T10:14:12Z", + "reconnect_count": 2, + "error": "" + }, + { + "channel": "telegram", + "status": "stale", + "last_event_at": "2026-08-06T09:30:00Z", + "reconnect_count": 0, + "error": "" + } + ] +} +``` + +Instance list responses include a compact `channel_health` summary so the +agent list can render a warning indicator without an extra request: + +```json +{ "overall": "unhealthy", "unhealthy_count": 1, "checked_at": "2026-08-06T10:15:00Z" } +``` + +## Data model + +`channel_health_statuses` stores the latest status per instance/channel +pair (one row per channel, overwritten on each check): + +| Field | Type | Description | +|---|---|---| +| `InstanceID` | uint | Instance the channel belongs to | +| `Channel` | string | Channel name (`slack`, `telegram`, …) | +| `Status` | string | One of the per-channel states above | +| `LastEventAt` | datetime | When the channel last received an event | +| `ReconnectCount` | int | Reconnects reported by the gateway | +| `Error` | string | Last error reported by the gateway, if any | +| `CheckedAt` | datetime | When this status was recorded | + +## UI + +- **Channel Health panel** on the Agent detail page (Settings tab): + per-channel status badges, last-event times, errors, and reconnect + counts. +- **Agent list**: agents whose overall status is unhealthy show a warning + indicator. + +## Known limitations + +- The gateway's channel state is in-memory OpenClaw state; it resets when + the gateway restarts, so last-event times and reconnect counts start + over. +- Staleness is inferred from event silence, so a genuinely quiet channel + (nobody messaging the agent for 30+ minutes) can be reported `stale` + even though it is fine. +- Monitoring only observes — it does not restart channels, restart + instances, or deliver alerts. + +## Future work + +- Auto-restart escalation for unhealthy channels. +- Alert delivery (notify operators when an agent goes unhealthy). +- Synthetic canary probes to distinguish quiet channels from stale ones. +- Consuming the gateway's push `health` broadcast instead of polling. diff --git a/website_docs/channel-health.mdx b/website_docs/channel-health.mdx new file mode 100644 index 00000000..c079b9ea --- /dev/null +++ b/website_docs/channel-health.mdx @@ -0,0 +1,75 @@ +--- +title: "Channel health" +description: "See whether your Agent's chat channels are connected and receiving events" +--- + +## Overview + +Claworc continuously checks whether each OpenClaw instance's chat channels — Slack, Telegram, Discord, +WhatsApp, and others — are connected and receiving events. About once a minute, the dashboard asks each +running OpenClaw instance for the state of its channels and shows the result, so you can spot a dropped +Slack connection or a misconfigured Telegram bot without opening logs. + +Monitoring is read-only: Claworc reports channel problems but does not restart anything on its own. + +## Where to find it + +Open an Agent and scroll the **Settings** tab to the **Channel Health** panel. For each channel it shows: + +- a **status badge** (see the table below) +- the time the channel **last received an event** +- the **last error** reported by the channel, if any +- how many times the channel has **reconnected** + +Agents with unhealthy channels also show a **warning indicator** in the Agents list, so you can spot +problems at a glance without opening each Agent. + +## What each status means + +| Status | Meaning | +|---|---| +| **Healthy** | The channel is connected and receiving events. | +| **Stale** | The channel's connection is up, but no events have arrived for over 30 minutes. Only applies to channels that hold a persistent connection (for example, Slack Socket Mode); webhook-based channels are never marked stale. | +| **Disconnected** | The channel is running, but its connection to the provider (Slack, Telegram, …) is down. | +| **Not running** | The channel is enabled in the OpenClaw instance's configuration but is not running. | +| **Disabled** | The channel is disabled or not configured on the OpenClaw instance. | + +The panel can also show a status for the Agent as a whole: + +| Status | Meaning | +|---|---| +| **Unreachable** | The OpenClaw instance is not responding — the OpenClaw process may be down. | +| **No channels** | The OpenClaw instance is reachable but has no chat channels configured. | +| **Unknown** | The Agent has not been checked yet (for example, it just started). | + +An Agent is considered **unhealthy** if any channel is Disconnected or Not running, and **degraded** +if any channel is Stale. + +## Fixing an unhealthy channel + +If a channel shows **Disconnected** or **Not running**: + +1. Check the channel's configuration on the OpenClaw instance — an expired or revoked token is the most + common cause. Update the token or credentials in the instance's OpenClaw config. +2. Restart the Agent from the dashboard. This restarts the OpenClaw instance and forces every channel to + reconnect. +3. If the problem persists, check the Agent's logs for errors from that channel. + +If a channel shows **Stale**, first consider whether it is simply quiet — a channel nobody has messaged +for half an hour is reported stale even when nothing is wrong. Send the Agent a test message on that +channel; if the message does not arrive, restart the Agent as above. + + + Statuses reset when the OpenClaw instance restarts, so last-event times and + reconnect counts start over after a restart. + + +## Configuration for operators + +Channel health monitoring is on by default. Operators can tune it with two environment variables on the +Claworc dashboard: + +| Env var | Default | Meaning | +|---|---|---| +| `CLAWORC_CHANNEL_HEALTH_ENABLED` | `true` | Enable or disable channel health monitoring | +| `CLAWORC_CHANNEL_HEALTH_INTERVAL` | `60s` | How often each Agent's channels are checked | From a9b5cb1d752c5eda5065266ad267aabe20147d3e Mon Sep 17 00:00:00 2001 From: Stan Misiurev Date: Thu, 6 Aug 2026 12:07:10 -0400 Subject: [PATCH 2/3] fix(frontend): rename channel health list badge to Warning Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01CfW5Dn8yEsULgbKzmcjbZ7 --- Makefile | 4 +++- .../frontend/src/common/components/ChannelHealthPanel.tsx | 2 +- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/Makefile b/Makefile index ec263fc5..7fb06da8 100644 --- a/Makefile +++ b/Makefile @@ -18,6 +18,8 @@ DASHBOARD_IMAGE := claworc/claworc TAG := latest PLATFORMS := linux/amd64,linux/arm64 NATIVE_ARCH := $(shell uname -m | sed 's/x86_64/amd64/') +# Go-installed dev tools (goreman, air) land here; not always on the user's PATH. +GOBIN_DIR := $(shell go env GOPATH)/bin CACHE_ARGS ?= @@ -177,7 +179,7 @@ dev: @echo "Control plane: http://localhost:8000" @echo "Frontend: http://localhost:5173" @echo "" - CLAWORC_AUTH_DISABLED=true CLAWORC_LLM_RESPONSE_LOG=$(CURDIR)/llm-responses.log CLAWORC_ALLOWED_HOST_MOUNTS=/tmp,~/ goreman -set-ports=false start + PATH="$(GOBIN_DIR):$$PATH" CLAWORC_AUTH_DISABLED=true CLAWORC_LLM_RESPONSE_LOG=$(CURDIR)/llm-responses.log CLAWORC_ALLOWED_HOST_MOUNTS=/tmp,~/ goreman -set-ports=false start ssh-integration-test: docker build -f agent/instance/Dockerfile -t claworc-agent:local agent/instance/ diff --git a/control-plane/frontend/src/common/components/ChannelHealthPanel.tsx b/control-plane/frontend/src/common/components/ChannelHealthPanel.tsx index b679de0d..4f8b09b1 100644 --- a/control-plane/frontend/src/common/components/ChannelHealthPanel.tsx +++ b/control-plane/frontend/src/common/components/ChannelHealthPanel.tsx @@ -181,7 +181,7 @@ export function ChannelHealthIndicator({ instance }: { instance: Instance }) { className="inline-flex items-center gap-1 px-2 py-0.5 rounded-full text-xs font-medium bg-red-100 text-red-800" > - Channels + Warning {tooltip} From fa8055907e39dbdd568f6c62fc2179bba8b3473d Mon Sep 17 00:00:00 2001 From: Stan Misiurev Date: Thu, 6 Aug 2026 13:35:53 -0400 Subject: [PATCH 3/3] feat(monitoring): escalate sustained channel failures with alerts and auto-restart The channel health monitor now feeds every snapshot to an escalator that tracks consecutive failing checks (overall unhealthy/unreachable) per instance and acts on sustained outages instead of only displaying them: - Webhook alert (channel_failure) after 3 consecutive failing checks, once per incident; recovery notification with outage duration when the incident closes. Payloads include a human-readable text field so Slack incoming webhooks work as-is; optional bearer token encrypted at rest. - Opt-in auto-restart (default off) after 5 consecutive failing checks, reusing the existing async restart flow, capped at 3 restarts per rolling hour per instance with a 10m post-restart cooldown. - channel_health_events audit table + GET /instances/{id}/channels/health/events. - Settings: Channel Health Alerts card (Misc tab) with webhook URL, token, toggles, and a Send Test button (POST /settings/channel-alerts/test). - Thresholds via CLAWORC_CHANNEL_HEALTH_{ALERT_THRESHOLD,RESTART_THRESHOLD, RESTART_MAX_PER_HOUR,RESTART_COOLDOWN}. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01CfW5Dn8yEsULgbKzmcjbZ7 --- .../frontend/src/app/pages/SettingsPage.tsx | 173 +++++++++++ .../frontend/src/common/api/settings.ts | 13 + .../frontend/src/common/types/settings.ts | 10 + .../channelhealth/channelhealth_test.go | 22 ++ .../internal/channelhealth/monitor.go | 17 ++ control-plane/internal/config/config.go | 11 + .../migrations/migration_00001_baseline.go | 1 + control-plane/internal/database/models.go | 1 + .../internal/database/models/models.go | 15 + .../internal/handlers/channel_alerts.go | 251 ++++++++++++++++ .../internal/handlers/channel_alerts_test.go | 137 +++++++++ .../internal/handlers/channel_escalation.go | 240 +++++++++++++++ .../handlers/channel_escalation_test.go | 283 ++++++++++++++++++ control-plane/internal/handlers/channels.go | 38 +++ control-plane/internal/handlers/settings.go | 25 +- control-plane/main.go | 11 + docs/channel-monitoring.md | 63 +++- 17 files changed, 1304 insertions(+), 7 deletions(-) create mode 100644 control-plane/internal/handlers/channel_alerts.go create mode 100644 control-plane/internal/handlers/channel_alerts_test.go create mode 100644 control-plane/internal/handlers/channel_escalation.go create mode 100644 control-plane/internal/handlers/channel_escalation_test.go diff --git a/control-plane/frontend/src/app/pages/SettingsPage.tsx b/control-plane/frontend/src/app/pages/SettingsPage.tsx index fb0d2b81..629dc61b 100644 --- a/control-plane/frontend/src/app/pages/SettingsPage.tsx +++ b/control-plane/frontend/src/app/pages/SettingsPage.tsx @@ -23,6 +23,7 @@ import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query"; import { useSettings, useUpdateSettings } from "@common/hooks/useSettings"; import { useProviders, useCatalogIconMap } from "@common/hooks/useProviders"; import { fetchSSHFingerprint, rotateSSHKey } from "@common/api/ssh"; +import { testChannelAlertWebhook } from "@common/api/settings"; import { syncAllProviders } from "@common/api/llm"; import { successToast, errorToast } from "@common/utils/toast"; import { validateResourceQuantities } from "@common/utils/resourceValidation"; @@ -800,6 +801,8 @@ function MiscTab({ )}
+ +

Anonymous Analytics

@@ -835,3 +838,173 @@ function MiscTab({

); } + +function ChannelAlertsCard({ settings }: { settings: Settings }) { + const updateMutation = useUpdateSettings(); + + const [url, setUrl] = useState(settings.channel_alert_webhook_url || ""); + const [alertsEnabled, setAlertsEnabled] = useState( + settings.channel_alerts_enabled !== "false", + ); + const [autoRestart, setAutoRestart] = useState( + settings.channel_auto_restart_enabled === "true", + ); + const [editingToken, setEditingToken] = useState(false); + const [showToken, setShowToken] = useState(false); + const [tokenValue, setTokenValue] = useState(""); + const [pendingToken, setPendingToken] = useState(null); + + const dirty = + url !== (settings.channel_alert_webhook_url || "") || + alertsEnabled !== (settings.channel_alerts_enabled !== "false") || + autoRestart !== (settings.channel_auto_restart_enabled === "true") || + pendingToken !== null; + + const save = () => { + const payload: SettingsUpdatePayload = { + channel_alert_webhook_url: url.trim(), + channel_alerts_enabled: alertsEnabled ? "true" : "false", + channel_auto_restart_enabled: autoRestart ? "true" : "false", + }; + if (pendingToken !== null) payload.channel_alert_webhook_token = pendingToken; + updateMutation.mutate(payload, { + onSuccess: () => { + setEditingToken(false); + setTokenValue(""); + setPendingToken(null); + }, + }); + }; + + const testMutation = useMutation({ + mutationFn: testChannelAlertWebhook, + onSuccess: (res) => { + if (res.status === "sent") { + successToast("Test alert delivered", `Webhook responded with HTTP ${res.http_status}`); + } else { + errorToast("Test alert failed", res.error || `HTTP ${res.http_status}`); + } + }, + onError: (err) => errorToast("Test alert failed", err), + }); + + return ( +
+

Channel Health Alerts

+

+ Get notified when an Agent's chat channels (Slack, Telegram, ...) stop responding. + Alerts are POSTed as JSON to the webhook URL below (Slack incoming webhooks work as-is). +

+
+
+ + setUrl(e.target.value)} + placeholder="https://hooks.slack.com/services/..." + className="w-full px-3 py-1.5 border border-gray-300 rounded-md text-sm focus:outline-none focus:ring-2 focus:ring-blue-500" + /> +
+
+ + {editingToken ? ( +
+
+ { + setTokenValue(e.target.value); + setPendingToken(e.target.value); + }} + className="w-full px-3 py-1.5 pr-10 border border-gray-300 rounded-md text-sm focus:outline-none focus:ring-2 focus:ring-blue-500" + placeholder="Sent as Authorization: Bearer " + /> + +
+ +
+ ) : ( +
+ + {pendingToken !== null + ? pendingToken + ? "****" + pendingToken.slice(-4) + : "(not set)" + : settings.channel_alert_webhook_token || "(not set)"} + + +
+ )} +
+ +
+ +

+ + Restarts the Agent's container after several consecutive failed checks + (max 3 restarts per hour). Thresholds are set via CLAWORC_CHANNEL_HEALTH_* env vars. +

+
+
+ + +
+
+
+ ); +} diff --git a/control-plane/frontend/src/common/api/settings.ts b/control-plane/frontend/src/common/api/settings.ts index dc999fc6..b2b5ae74 100644 --- a/control-plane/frontend/src/common/api/settings.ts +++ b/control-plane/frontend/src/common/api/settings.ts @@ -12,3 +12,16 @@ export async function updateSettings( const { data } = await client.put("/settings", payload); return data; } + +export interface ChannelAlertTestResult { + status: "sent" | "failed"; + http_status?: number; + error?: string; +} + +export async function testChannelAlertWebhook(): Promise { + const { data } = await client.post( + "/settings/channel-alerts/test", + ); + return data; +} diff --git a/control-plane/frontend/src/common/types/settings.ts b/control-plane/frontend/src/common/types/settings.ts index cbdc2ee0..26fe7f97 100644 --- a/control-plane/frontend/src/common/types/settings.ts +++ b/control-plane/frontend/src/common/types/settings.ts @@ -33,6 +33,12 @@ export interface Settings { analytics_consent: "unset" | "opt_in" | "opt_out"; /** Random 32-char hex ID reported alongside anonymous events. Read-only. */ installation_id: string; + /** Channel alert delivery. Booleans stored as "true"/"false" ("" = default). */ + channel_alerts_enabled: string; + channel_auto_restart_enabled: string; + channel_alert_webhook_url: string; + /** Masked (e.g. "****abcd") — write-only via update. */ + channel_alert_webhook_token: string; /** * Only populated on the PUT response when env vars changed: the set of * running instances the backend kicked a restart on to apply the change. @@ -64,6 +70,10 @@ export interface SettingsUpdatePayload { default_affinity?: string; default_service_account_annotations?: Record; default_ports?: import("./instance").PortSpec[]; + channel_alerts_enabled?: string; + channel_auto_restart_enabled?: string; + channel_alert_webhook_url?: string; + channel_alert_webhook_token?: string; } // Keep backward compat alias diff --git a/control-plane/internal/channelhealth/channelhealth_test.go b/control-plane/internal/channelhealth/channelhealth_test.go index eedbd1dd..2619c32e 100644 --- a/control-plane/internal/channelhealth/channelhealth_test.go +++ b/control-plane/internal/channelhealth/channelhealth_test.go @@ -281,3 +281,25 @@ func TestDeriveOverall(t *testing.T) { }) } } + +func TestListenerFiresOnEveryStore(t *testing.T) { + m := New(nil, time.Minute) + var got []string + m.SetListener(func(snap Snapshot) { got = append(got, snap.Overall) }) + + // Same overall twice: the listener must fire both times even though + // the transition-logging path early-returns on no-change. + m.store(Snapshot{InstanceID: 1, Overall: OverallUnhealthy}) + m.store(Snapshot{InstanceID: 1, Overall: OverallUnhealthy}) + m.store(Snapshot{InstanceID: 1, Overall: OverallHealthy}) + + want := []string{OverallUnhealthy, OverallUnhealthy, OverallHealthy} + if len(got) != len(want) { + t.Fatalf("expected %d listener calls, got %d", len(want), len(got)) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("call %d: expected %q, got %q", i, want[i], got[i]) + } + } +} diff --git a/control-plane/internal/channelhealth/monitor.go b/control-plane/internal/channelhealth/monitor.go index 39c40c2b..125e3dc5 100644 --- a/control-plane/internal/channelhealth/monitor.go +++ b/control-plane/internal/channelhealth/monitor.go @@ -25,12 +25,19 @@ const ( gatewayTunnelLabel = "Gateway" ) +// Listener receives every stored snapshot, including ones whose overall +// status did not change — consumers that count consecutive results depend +// on non-transition snapshots too. Listeners run synchronously on the +// check goroutine and must not block. +type Listener func(snap Snapshot) + // Monitor periodically polls the OpenClaw gateway of every running // instance for channel health, persists the results, and keeps an // in-memory snapshot per instance for cheap reads by handlers. type Monitor struct { tunnels *sshproxy.TunnelManager interval time.Duration + listener Listener mu sync.RWMutex snapshots map[uint]Snapshot @@ -50,6 +57,12 @@ func New(tunnels *sshproxy.TunnelManager, interval time.Duration) *Monitor { } } +// SetListener registers the snapshot listener. Must be called before +// Start; the field is not synchronized. +func (m *Monitor) SetListener(fn Listener) { + m.listener = fn +} + // Start launches the background polling loop. It returns immediately; the // goroutine exits when ctx is canceled. func (m *Monitor) Start(ctx context.Context) { @@ -190,6 +203,10 @@ func (m *Monitor) store(snap Snapshot) { m.snapshots[snap.InstanceID] = snap m.mu.Unlock() + if m.listener != nil { + m.listener(snap) + } + prevOverall := OverallUnknown if had { prevOverall = prev.Overall diff --git a/control-plane/internal/config/config.go b/control-plane/internal/config/config.go index 0b633245..55950a18 100644 --- a/control-plane/internal/config/config.go +++ b/control-plane/internal/config/config.go @@ -54,6 +54,17 @@ type Settings struct { // makes the channel health endpoints report "disabled". ChannelHealthEnabled bool `envconfig:"CHANNEL_HEALTH_ENABLED" default:"true"` ChannelHealthInterval time.Duration `envconfig:"CHANNEL_HEALTH_INTERVAL" default:"60s"` + + // Channel health escalation thresholds, counted in consecutive failing + // checks (overall unhealthy/unreachable). Alerts fire once per incident + // at the alert threshold; auto-restart (opt-in via the + // channel_auto_restart_enabled setting) fires at the restart threshold, + // capped per rolling hour and followed by a cooldown that covers the + // container rebuild and channel reconnect window. + ChannelHealthAlertThreshold int `envconfig:"CHANNEL_HEALTH_ALERT_THRESHOLD" default:"3"` + ChannelHealthRestartThreshold int `envconfig:"CHANNEL_HEALTH_RESTART_THRESHOLD" default:"5"` + ChannelHealthRestartMaxPerHour int `envconfig:"CHANNEL_HEALTH_RESTART_MAX_PER_HOUR" default:"3"` + ChannelHealthRestartCooldown time.Duration `envconfig:"CHANNEL_HEALTH_RESTART_COOLDOWN" default:"10m"` } var Cfg Settings diff --git a/control-plane/internal/database/migrations/migration_00001_baseline.go b/control-plane/internal/database/migrations/migration_00001_baseline.go index 4c9a0db0..308d7e17 100644 --- a/control-plane/internal/database/migrations/migration_00001_baseline.go +++ b/control-plane/internal/database/migrations/migration_00001_baseline.go @@ -64,5 +64,6 @@ func AutoMigrateAll(gdb interface { &models.WebhookApiKey{}, &models.WebhookLog{}, &models.ChannelHealthStatus{}, + &models.ChannelHealthEvent{}, ) } diff --git a/control-plane/internal/database/models.go b/control-plane/internal/database/models.go index ed475d62..12218e8f 100644 --- a/control-plane/internal/database/models.go +++ b/control-plane/internal/database/models.go @@ -36,6 +36,7 @@ type ( WebhookApiKey = models.WebhookApiKey WebhookLog = models.WebhookLog ChannelHealthStatus = models.ChannelHealthStatus + ChannelHealthEvent = models.ChannelHealthEvent ) // Helper re-exports keep `database.ParseTeamIDs(...)` etc. working for diff --git a/control-plane/internal/database/models/models.go b/control-plane/internal/database/models/models.go index 5c35d57f..f0841efe 100644 --- a/control-plane/internal/database/models/models.go +++ b/control-plane/internal/database/models/models.go @@ -471,6 +471,21 @@ type ChannelHealthStatus struct { UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"` } +// ChannelHealthEvent is a durable audit record of channel-health +// escalation actions (alerts, auto-restarts, recoveries). It doubles as +// outage history: a failure_detected/recovered pair brackets an incident. +type ChannelHealthEvent struct { + ID uint `gorm:"primaryKey;autoIncrement" json:"id"` + InstanceID uint `gorm:"not null;index" json:"instance_id"` + Type string `gorm:"not null" json:"type"` // failure_detected|auto_restart|restart_limit_reached|recovered|webhook_test + Overall string `gorm:"default:''" json:"overall"` + // Detail is a JSON blob with incident context (failing channels, + // consecutive check count, outage duration). + Detail string `gorm:"type:text;default:''" json:"detail"` + WebhookStatus string `gorm:"default:''" json:"webhook_status"` // sent|failed|skipped + CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` +} + // UserSSHKey is a public key a user authenticates with against the inbound // SSH gateway. The private key is never stored — it is generated on demand // and handed to the user exactly once (or the user uploads their own pubkey). diff --git a/control-plane/internal/handlers/channel_alerts.go b/control-plane/internal/handlers/channel_alerts.go new file mode 100644 index 00000000..2718e4ce --- /dev/null +++ b/control-plane/internal/handlers/channel_alerts.go @@ -0,0 +1,251 @@ +package handlers + +import ( + "bytes" + "encoding/json" + "fmt" + "log" + "net/http" + "time" + + "github.com/gluk-w/claworc/control-plane/internal/channelhealth" + "github.com/gluk-w/claworc/control-plane/internal/database" + "github.com/gluk-w/claworc/control-plane/internal/utils" +) + +// Settings keys for the channel alert webhook. The URL is a plain setting; +// the bearer token is encrypted at rest like brave_api_key. +const ( + settingChannelAlertsEnabled = "channel_alerts_enabled" + settingChannelAlertWebhookURL = "channel_alert_webhook_url" + settingChannelAlertWebhookToken = "channel_alert_webhook_token" +) + +// Webhook delivery outcomes recorded on ChannelHealthEvent rows. +const ( + webhookStatusSent = "sent" + webhookStatusFailed = "failed" + webhookStatusSkipped = "skipped" +) + +var channelAlertClient = &http.Client{Timeout: 10 * time.Second} + +// channelAlertRetryDelay is overridable in tests. +var channelAlertRetryDelay = 5 * time.Second + +// AlertInstance identifies the instance an alert is about. +type AlertInstance struct { + ID uint `json:"id"` + Name string `json:"name"` + DisplayName string `json:"display_name"` +} + +// AlertChannel is one non-healthy channel included in an alert. +type AlertChannel struct { + Channel string `json:"channel"` + AccountID string `json:"account_id"` + Status string `json:"status"` + LastError string `json:"last_error,omitempty"` +} + +// ChannelAlertPayload is the JSON body POSTed to the configured channel +// alert webhook. The Text field is a self-contained human-readable summary +// so bare Slack/Discord incoming-webhook style receivers are useful as-is. +type ChannelAlertPayload struct { + Event string `json:"event"` // channel_failure|auto_restart|restart_limit_reached|recovery|test + Text string `json:"text"` + Timestamp time.Time `json:"timestamp"` + Instance AlertInstance `json:"instance"` + Overall string `json:"overall,omitempty"` + ConsecutiveFailures int `json:"consecutive_failures,omitempty"` + FailingSince *time.Time `json:"failing_since,omitempty"` + DurationSeconds int `json:"duration_seconds,omitempty"` + Channels []AlertChannel `json:"channels,omitempty"` +} + +// buildChannelAlertPayload assembles the webhook payload for an escalation +// event from the health snapshot that triggered it. +func buildChannelAlertPayload(snap channelhealth.Snapshot, eventType string, extra map[string]any) ChannelAlertPayload { + var inst database.Instance + _ = database.DB.First(&inst, snap.InstanceID).Error + + p := ChannelAlertPayload{ + Timestamp: time.Now().UTC(), + Overall: snap.Overall, + Instance: AlertInstance{ + ID: snap.InstanceID, + Name: inst.Name, + DisplayName: inst.DisplayName, + }, + } + for _, ch := range snap.Channels { + if ch.Status == channelhealth.StatusHealthy || ch.Status == channelhealth.StatusDisabled { + continue + } + p.Channels = append(p.Channels, AlertChannel{ + Channel: ch.Channel, + AccountID: ch.AccountID, + Status: ch.Status, + LastError: ch.LastError, + }) + } + if v, ok := extra["consecutive_failures"].(int); ok { + p.ConsecutiveFailures = v + } + if v, ok := extra["duration_seconds"].(int); ok { + p.DurationSeconds = v + } + if v, ok := extra["failing_since"].(time.Time); ok && !v.IsZero() { + t := v + p.FailingSince = &t + } + + name := inst.DisplayName + if name == "" { + name = fmt.Sprintf("#%d", snap.InstanceID) + } + chansText := "" + for i, ch := range p.Channels { + if i > 0 { + chansText += ", " + } + chansText += fmt.Sprintf("%s/%s: %s", ch.Channel, ch.AccountID, ch.Status) + } + if chansText != "" { + chansText = " (" + chansText + ")" + } + + switch eventType { + case eventFailureDetected: + p.Event = "channel_failure" + p.Text = fmt.Sprintf("Claworc: agent %q channels %s for %d consecutive checks%s", + name, snap.Overall, p.ConsecutiveFailures, chansText) + case eventAutoRestart: + p.Event = "auto_restart" + p.Text = fmt.Sprintf("Claworc: auto-restarting agent %q — channels %s for %d consecutive checks%s", + name, snap.Overall, p.ConsecutiveFailures, chansText) + case eventRestartLimitReached: + p.Event = "restart_limit_reached" + p.Text = fmt.Sprintf("Claworc: agent %q still %s but the auto-restart limit was reached; manual intervention needed%s", + name, snap.Overall, chansText) + case eventRecovered: + p.Event = "recovery" + p.Text = fmt.Sprintf("Claworc: agent %q channel health recovered after %s", + name, (time.Duration(p.DurationSeconds) * time.Second).String()) + default: + p.Event = eventType + p.Text = fmt.Sprintf("Claworc: agent %q channel health event %q", name, eventType) + } + return p +} + +// channelAlertConfig reads the alert delivery settings. Returns ok=false +// when alerts are disabled or no URL is configured. +func channelAlertConfig() (url, token string, ok bool) { + if enabled, err := database.GetSetting(settingChannelAlertsEnabled); err == nil && enabled == "false" { + return "", "", false + } + url, err := database.GetSetting(settingChannelAlertWebhookURL) + if err != nil || url == "" { + return "", "", false + } + if enc, err := database.GetSetting(settingChannelAlertWebhookToken); err == nil && enc != "" { + if tok, err := utils.Decrypt(enc); err == nil { + token = tok + } + } + return url, token, true +} + +// sendChannelAlert delivers the payload to the configured webhook with one +// retry on network error or 5xx. Returns the delivery outcome for the +// audit row. Callers must not invoke this on the monitor goroutine. +func sendChannelAlert(p ChannelAlertPayload) string { + url, token, ok := channelAlertConfig() + if !ok { + return webhookStatusSkipped + } + status, _, err := postChannelAlert(url, token, p) + if err != nil || status >= 500 { + time.Sleep(channelAlertRetryDelay) + status, _, err = postChannelAlert(url, token, p) + } + if err != nil { + log.Printf("[channelalert] delivery failed: %v", err) + return webhookStatusFailed + } + if status >= 300 { + log.Printf("[channelalert] delivery failed: HTTP %d", status) + return webhookStatusFailed + } + return webhookStatusSent +} + +func postChannelAlert(url, token string, p ChannelAlertPayload) (int, string, error) { + body, err := json.Marshal(p) + if err != nil { + return 0, "", err + } + req, err := http.NewRequest(http.MethodPost, url, bytes.NewReader(body)) + if err != nil { + return 0, "", err + } + req.Header.Set("Content-Type", "application/json") + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + resp, err := channelAlertClient.Do(req) + if err != nil { + return 0, "", err + } + defer resp.Body.Close() + return resp.StatusCode, resp.Status, nil +} + +// TestChannelAlertWebhook sends a synchronous test alert to the configured +// webhook so admins can verify delivery from the Settings page. +// POST /api/v1/settings/channel-alerts/test (admin only). +func TestChannelAlertWebhook(w http.ResponseWriter, r *http.Request) { + url, token, ok := channelAlertConfig() + if !ok { + writeError(w, http.StatusBadRequest, "Channel alerts are disabled or no webhook URL is configured") + return + } + p := ChannelAlertPayload{ + Event: "test", + Text: "Claworc: test alert — channel alert webhook is configured correctly", + Timestamp: time.Now().UTC(), + } + status, statusText, err := postChannelAlert(url, token, p) + + ev := database.ChannelHealthEvent{ + Type: "webhook_test", + WebhookStatus: webhookStatusSent, + } + if err != nil || status >= 300 { + ev.WebhookStatus = webhookStatusFailed + } + if dbErr := database.DB.Create(&ev).Error; dbErr != nil { + log.Printf("[channelalert] record test event: %v", dbErr) + } + + if err != nil { + writeJSON(w, http.StatusBadGateway, map[string]any{ + "status": "failed", + "error": err.Error(), + }) + return + } + if status >= 300 { + writeJSON(w, http.StatusBadGateway, map[string]any{ + "status": "failed", + "http_status": status, + "error": statusText, + }) + return + } + writeJSON(w, http.StatusOK, map[string]any{ + "status": "sent", + "http_status": status, + }) +} diff --git a/control-plane/internal/handlers/channel_alerts_test.go b/control-plane/internal/handlers/channel_alerts_test.go new file mode 100644 index 00000000..4de5d0b9 --- /dev/null +++ b/control-plane/internal/handlers/channel_alerts_test.go @@ -0,0 +1,137 @@ +package handlers + +import ( + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "github.com/gluk-w/claworc/control-plane/internal/channelhealth" + "github.com/gluk-w/claworc/control-plane/internal/database" + "github.com/gluk-w/claworc/control-plane/internal/utils" +) + +func configureAlertWebhook(t *testing.T, url, token string) { + t.Helper() + if err := database.SetSetting(settingChannelAlertWebhookURL, url); err != nil { + t.Fatalf("set url: %v", err) + } + if token != "" { + enc, err := utils.Encrypt(token) + if err != nil { + t.Fatalf("encrypt token: %v", err) + } + if err := database.SetSetting(settingChannelAlertWebhookToken, enc); err != nil { + t.Fatalf("set token: %v", err) + } + } +} + +func TestSendChannelAlert_DeliversPayloadWithBearer(t *testing.T) { + setupHandlersTestDB(t) + + var gotBody []byte + var gotAuth string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody, _ = io.ReadAll(r.Body) + gotAuth = r.Header.Get("Authorization") + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + configureAlertWebhook(t, srv.URL, "sekret") + + p := buildChannelAlertPayload(channelhealth.Snapshot{ + InstanceID: 42, + Overall: channelhealth.OverallUnhealthy, + Channels: []channelhealth.ChannelState{ + {Channel: "slack", AccountID: "default", Status: channelhealth.StatusDisconnected, LastError: "socket closed"}, + {Channel: "telegram", AccountID: "default", Status: channelhealth.StatusHealthy}, + }, + }, eventFailureDetected, map[string]any{"consecutive_failures": 3}) + + if got := sendChannelAlert(p); got != webhookStatusSent { + t.Fatalf("expected sent, got %q", got) + } + if gotAuth != "Bearer sekret" { + t.Fatalf("expected bearer header, got %q", gotAuth) + } + + var decoded ChannelAlertPayload + if err := json.Unmarshal(gotBody, &decoded); err != nil { + t.Fatalf("payload not JSON: %v", err) + } + if decoded.Event != "channel_failure" || decoded.Text == "" { + t.Fatalf("unexpected payload: %+v", decoded) + } + if len(decoded.Channels) != 1 || decoded.Channels[0].Channel != "slack" { + t.Fatalf("expected only non-healthy channels, got %+v", decoded.Channels) + } + if decoded.ConsecutiveFailures != 3 { + t.Fatalf("expected consecutive_failures=3, got %d", decoded.ConsecutiveFailures) + } +} + +func TestSendChannelAlert_RetriesOn5xx(t *testing.T) { + setupHandlersTestDB(t) + + old := channelAlertRetryDelay + channelAlertRetryDelay = time.Millisecond + t.Cleanup(func() { channelAlertRetryDelay = old }) + + var calls atomic.Int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if calls.Add(1) == 1 { + w.WriteHeader(http.StatusInternalServerError) + return + } + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + configureAlertWebhook(t, srv.URL, "") + + if got := sendChannelAlert(ChannelAlertPayload{Event: "test"}); got != webhookStatusSent { + t.Fatalf("expected sent after retry, got %q", got) + } + if calls.Load() != 2 { + t.Fatalf("expected 2 attempts, got %d", calls.Load()) + } +} + +func TestSendChannelAlert_FailsAfterRetry(t *testing.T) { + setupHandlersTestDB(t) + + old := channelAlertRetryDelay + channelAlertRetryDelay = time.Millisecond + t.Cleanup(func() { channelAlertRetryDelay = old }) + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + defer srv.Close() + configureAlertWebhook(t, srv.URL, "") + + if got := sendChannelAlert(ChannelAlertPayload{Event: "test"}); got != webhookStatusFailed { + t.Fatalf("expected failed, got %q", got) + } +} + +func TestSendChannelAlert_SkippedWhenUnconfigured(t *testing.T) { + setupHandlersTestDB(t) + if got := sendChannelAlert(ChannelAlertPayload{Event: "test"}); got != webhookStatusSkipped { + t.Fatalf("expected skipped without URL, got %q", got) + } +} + +func TestSendChannelAlert_SkippedWhenDisabled(t *testing.T) { + setupHandlersTestDB(t) + configureAlertWebhook(t, "http://127.0.0.1:1/never", "") + if err := database.SetSetting(settingChannelAlertsEnabled, "false"); err != nil { + t.Fatalf("set setting: %v", err) + } + if got := sendChannelAlert(ChannelAlertPayload{Event: "test"}); got != webhookStatusSkipped { + t.Fatalf("expected skipped when disabled, got %q", got) + } +} diff --git a/control-plane/internal/handlers/channel_escalation.go b/control-plane/internal/handlers/channel_escalation.go new file mode 100644 index 00000000..e1159898 --- /dev/null +++ b/control-plane/internal/handlers/channel_escalation.go @@ -0,0 +1,240 @@ +package handlers + +import ( + "encoding/json" + "fmt" + "log" + "sync" + "time" + + "github.com/gluk-w/claworc/control-plane/internal/channelhealth" + "github.com/gluk-w/claworc/control-plane/internal/database" +) + +// settingChannelAutoRestartEnabled is the DB settings key gating automatic +// restarts. Auto-restart is opt-in: absent or non-"true" means disabled. +const settingChannelAutoRestartEnabled = "channel_auto_restart_enabled" + +// Channel health event types persisted to channel_health_events. +const ( + eventFailureDetected = "failure_detected" + eventAutoRestart = "auto_restart" + eventRestartLimitReached = "restart_limit_reached" + eventRecovered = "recovered" +) + +// restartWindow is the rolling window the auto-restart circuit breaker +// counts restarts in. +const restartWindow = time.Hour + +// ChannelEscalatorConfig holds the escalation thresholds (see +// CLAWORC_CHANNEL_HEALTH_* env vars). +type ChannelEscalatorConfig struct { + AlertThreshold int + RestartThreshold int + MaxRestartsPerHour int + RestartCooldown time.Duration +} + +// escalationState is the in-memory incident state for one instance. It is +// not persisted: after a control-plane restart an ongoing outage re-counts +// from zero. +type escalationState struct { + consecutiveFails int + failingSince time.Time + alertSent bool + breakerAlerted bool + cooldownUntil time.Time + // restartTimes is the rolling-window restart log for the circuit + // breaker. Deliberately preserved across incident resets so a + // restart -> briefly-healthy -> fail loop cannot restart forever. + restartTimes []time.Time +} + +// ChannelEscalator turns channel health snapshots into alerts and (opt-in) +// automatic instance restarts. It is registered as the channelhealth +// Monitor's listener; OnSnapshot does only in-memory bookkeeping +// synchronously and dispatches all I/O (DB writes, webhook, restart) in +// goroutines so the monitor loop never blocks. +type ChannelEscalator struct { + cfg ChannelEscalatorConfig + + // Injectable for tests. + now func() time.Time + restart func(instanceID uint, title, message string) + notify func(p ChannelAlertPayload) string + recordEvent func(ev database.ChannelHealthEvent) + // dispatch runs slow work off the monitor goroutine (tests run it + // inline for determinism). + dispatch func(fn func()) + + mu sync.Mutex + states map[uint]*escalationState +} + +// NewChannelEscalator builds an escalator with production dependencies. +func NewChannelEscalator(cfg ChannelEscalatorConfig) *ChannelEscalator { + e := &ChannelEscalator{ + cfg: cfg, + now: time.Now, + notify: sendChannelAlert, + states: make(map[uint]*escalationState), + dispatch: func(fn func()) { go fn() }, + } + e.restart = e.restartInstance + e.recordEvent = func(ev database.ChannelHealthEvent) { + if err := database.DB.Create(&ev).Error; err != nil { + log.Printf("[channelhealth] record event: %v", err) + } + } + return e +} + +// OnSnapshot is the channelhealth.Listener. It receives every stored +// snapshot, including ones whose overall status did not change. +func (e *ChannelEscalator) OnSnapshot(snap channelhealth.Snapshot) { + failing := snap.Overall == channelhealth.OverallUnhealthy || snap.Overall == channelhealth.OverallUnreachable + recovered := snap.Overall == channelhealth.OverallHealthy || snap.Overall == channelhealth.OverallNoChannels + // degraded/unknown hold: neither count nor reset an open incident. + if !failing && !recovered { + return + } + + e.mu.Lock() + defer e.mu.Unlock() + st := e.states[snap.InstanceID] + if st == nil { + st = &escalationState{} + e.states[snap.InstanceID] = st + } + now := e.now() + + if recovered { + if st.alertSent { + duration := now.Sub(st.failingSince) + e.dispatch(func() { + e.emit(snap, eventRecovered, map[string]any{ + "duration_seconds": int(duration.Seconds()), + }) + }) + } + st.consecutiveFails = 0 + st.failingSince = time.Time{} + st.alertSent = false + st.breakerAlerted = false + st.cooldownUntil = time.Time{} + return + } + + // Failing snapshot. + if now.Before(st.cooldownUntil) { + return + } + st.consecutiveFails++ + if st.consecutiveFails == 1 { + st.failingSince = snap.CheckedAt + if st.failingSince.IsZero() { + st.failingSince = now + } + } + + if st.consecutiveFails >= e.cfg.AlertThreshold && !st.alertSent { + st.alertSent = true + fails := st.consecutiveFails + since := st.failingSince + e.dispatch(func() { + e.emit(snap, eventFailureDetected, map[string]any{ + "consecutive_failures": fails, + "failing_since": since, + }) + }) + } + + if st.consecutiveFails < e.cfg.RestartThreshold { + return + } + if !autoRestartEnabled() { + return + } + + // Circuit breaker: cap restarts per instance per rolling hour. + cutoff := now.Add(-restartWindow) + recent := st.restartTimes[:0] + for _, t := range st.restartTimes { + if t.After(cutoff) { + recent = append(recent, t) + } + } + st.restartTimes = recent + if len(st.restartTimes) >= e.cfg.MaxRestartsPerHour { + if !st.breakerAlerted { + st.breakerAlerted = true + restarts := len(st.restartTimes) + e.dispatch(func() { + e.emit(snap, eventRestartLimitReached, map[string]any{ + "restarts_last_hour": restarts, + }) + }) + } + return + } + + st.restartTimes = append(st.restartTimes, now) + st.cooldownUntil = now.Add(e.cfg.RestartCooldown) + fails := st.consecutiveFails + since := st.failingSince + e.dispatch(func() { + e.restart(snap.InstanceID, + "Auto-restarting agent with unhealthy channels", + fmt.Sprintf("Channel health %s for %d consecutive checks", snap.Overall, fails)) + e.emit(snap, eventAutoRestart, map[string]any{ + "consecutive_failures": fails, + "failing_since": since, + }) + }) +} + +// emit persists an audit event and sends the webhook alert for it. Runs +// off the monitor goroutine. +func (e *ChannelEscalator) emit(snap channelhealth.Snapshot, eventType string, extra map[string]any) { + payload := buildChannelAlertPayload(snap, eventType, extra) + status := e.notify(payload) + + detail := map[string]any{} + for k, v := range extra { + detail[k] = v + } + if len(payload.Channels) > 0 { + detail["channels"] = payload.Channels + } + detailJSON, _ := json.Marshal(detail) + e.recordEvent(database.ChannelHealthEvent{ + InstanceID: snap.InstanceID, + Type: eventType, + Overall: snap.Overall, + Detail: string(detailJSON), + WebhookStatus: status, + }) +} + +// restartInstance is the production restart dependency: it re-fetches a +// fresh instance row (the snapshot may be up to one interval old) and +// reuses the shared async restart flow, which no-ops unless the instance +// is still running. +func (e *ChannelEscalator) restartInstance(instanceID uint, title, message string) { + var inst database.Instance + if err := database.DB.First(&inst, instanceID).Error; err != nil { + log.Printf("[channelhealth] auto-restart: load instance %d: %v", instanceID, err) + return + } + log.Printf("[channelhealth] auto-restarting instance %d (%s): %s", inst.ID, inst.DisplayName, message) + restartInstanceAsyncWithToast(inst, 0, title, message) +} + +// autoRestartEnabled reads the opt-in toggle from the settings table. +// Settings are intentionally not cached (matches the rest of the settings +// surface), so flipping the toggle takes effect on the next check. +func autoRestartEnabled() bool { + val, err := database.GetSetting(settingChannelAutoRestartEnabled) + return err == nil && val == "true" +} diff --git a/control-plane/internal/handlers/channel_escalation_test.go b/control-plane/internal/handlers/channel_escalation_test.go new file mode 100644 index 00000000..f54e70b0 --- /dev/null +++ b/control-plane/internal/handlers/channel_escalation_test.go @@ -0,0 +1,283 @@ +package handlers + +import ( + "sync" + "testing" + "time" + + "github.com/gluk-w/claworc/control-plane/internal/channelhealth" + "github.com/gluk-w/claworc/control-plane/internal/database" +) + +type escalatorFixture struct { + esc *ChannelEscalator + clock time.Time + mu sync.Mutex + restarts []uint + notified []ChannelAlertPayload + events []database.ChannelHealthEvent +} + +func newEscalatorFixture(t *testing.T) *escalatorFixture { + t.Helper() + setupHandlersTestDB(t) + + f := &escalatorFixture{clock: time.Date(2026, 8, 6, 12, 0, 0, 0, time.UTC)} + f.esc = NewChannelEscalator(ChannelEscalatorConfig{ + AlertThreshold: 3, + RestartThreshold: 5, + MaxRestartsPerHour: 3, + RestartCooldown: 10 * time.Minute, + }) + f.esc.now = func() time.Time { return f.clock } + f.esc.dispatch = func(fn func()) { fn() } + f.esc.restart = func(id uint, title, message string) { + f.mu.Lock() + defer f.mu.Unlock() + f.restarts = append(f.restarts, id) + } + f.esc.notify = func(p ChannelAlertPayload) string { + f.mu.Lock() + defer f.mu.Unlock() + f.notified = append(f.notified, p) + return webhookStatusSent + } + f.esc.recordEvent = func(ev database.ChannelHealthEvent) { + f.mu.Lock() + defer f.mu.Unlock() + f.events = append(f.events, ev) + } + return f +} + +func (f *escalatorFixture) advance(d time.Duration) { f.clock = f.clock.Add(d) } + +func (f *escalatorFixture) snap(overall string) channelhealth.Snapshot { + return channelhealth.Snapshot{ + InstanceID: 1, + Overall: overall, + CheckedAt: f.clock, + Channels: []channelhealth.ChannelState{{ + Channel: "slack", AccountID: "default", + Status: channelhealth.StatusDisconnected, + }}, + } +} + +// tick feeds one failing snapshot and advances the clock one interval. +func (f *escalatorFixture) tick(overall string) { + f.esc.OnSnapshot(f.snap(overall)) + f.advance(time.Minute) +} + +func (f *escalatorFixture) eventTypes() []string { + f.mu.Lock() + defer f.mu.Unlock() + out := make([]string, len(f.events)) + for i, e := range f.events { + out[i] = e.Type + } + return out +} + +func enableAutoRestart(t *testing.T) { + t.Helper() + if err := database.SetSetting(settingChannelAutoRestartEnabled, "true"); err != nil { + t.Fatalf("set setting: %v", err) + } +} + +func TestEscalator_NoActionBelowAlertThreshold(t *testing.T) { + f := newEscalatorFixture(t) + f.tick(channelhealth.OverallUnhealthy) + f.tick(channelhealth.OverallUnhealthy) + if len(f.events) != 0 || len(f.notified) != 0 { + t.Fatalf("expected no actions below threshold, got events=%v", f.eventTypes()) + } +} + +func TestEscalator_AlertOncePerIncident(t *testing.T) { + f := newEscalatorFixture(t) + for i := 0; i < 4; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + if got := f.eventTypes(); len(got) != 1 || got[0] != eventFailureDetected { + t.Fatalf("expected one failure_detected, got %v", got) + } + if f.notified[0].Event != "channel_failure" { + t.Fatalf("expected channel_failure payload, got %q", f.notified[0].Event) + } + if f.notified[0].ConsecutiveFailures != 3 { + t.Fatalf("expected 3 consecutive failures in payload, got %d", f.notified[0].ConsecutiveFailures) + } + if f.notified[0].FailingSince == nil { + t.Fatal("expected failing_since to be set") + } +} + +func TestEscalator_NoRestartWhenToggleOff(t *testing.T) { + f := newEscalatorFixture(t) + for i := 0; i < 8; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + if len(f.restarts) != 0 { + t.Fatalf("expected no restarts with toggle off, got %d", len(f.restarts)) + } +} + +func TestEscalator_RestartAtThresholdThenCooldown(t *testing.T) { + f := newEscalatorFixture(t) + enableAutoRestart(t) + for i := 0; i < 5; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + if len(f.restarts) != 1 { + t.Fatalf("expected exactly one restart at threshold, got %d", len(f.restarts)) + } + if got := f.eventTypes(); len(got) != 2 || got[1] != eventAutoRestart { + t.Fatalf("expected [failure_detected auto_restart], got %v", got) + } + // Failing checks during the 10m cooldown are ignored entirely. + for i := 0; i < 9; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + if len(f.restarts) != 1 { + t.Fatalf("cooldown violated: got %d restarts", len(f.restarts)) + } + // After cooldown the counter restarts from zero: 5 more failing checks + // trigger the second restart. + for i := 0; i < 5; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + if len(f.restarts) != 2 { + t.Fatalf("expected second restart after cooldown + threshold, got %d", len(f.restarts)) + } +} + +func TestEscalator_CircuitBreaker(t *testing.T) { + f := newEscalatorFixture(t) + enableAutoRestart(t) + // Drive three restarts (threshold 5 fails + 10m cooldown between). + for r := 0; r < 3; r++ { + for i := 0; i < 5; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + f.advance(10 * time.Minute) + } + if len(f.restarts) != 3 { + t.Fatalf("expected 3 restarts before breaker, got %d", len(f.restarts)) + } + // Fourth attempt within the hour trips the breaker instead. + for i := 0; i < 5; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + if len(f.restarts) != 3 { + t.Fatalf("breaker violated: got %d restarts", len(f.restarts)) + } + types := f.eventTypes() + if types[len(types)-1] != eventRestartLimitReached { + t.Fatalf("expected restart_limit_reached, got %v", types) + } + // Breaker alert fires once per incident even as failures continue. + for i := 0; i < 5; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + count := 0 + for _, tp := range f.eventTypes() { + if tp == eventRestartLimitReached { + count++ + } + } + if count != 1 { + t.Fatalf("expected one restart_limit_reached, got %d", count) + } +} + +func TestEscalator_RecoveryResetsIncidentKeepsRestartLog(t *testing.T) { + f := newEscalatorFixture(t) + enableAutoRestart(t) + for r := 0; r < 3; r++ { + for i := 0; i < 5; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + f.advance(10 * time.Minute) + } + f.tick(channelhealth.OverallHealthy) + types := f.eventTypes() + if types[len(types)-1] != eventRecovered { + t.Fatalf("expected recovered event, got %v", types) + } + if f.notified[len(f.notified)-1].DurationSeconds <= 0 { + t.Fatal("expected positive outage duration in recovery payload") + } + // New incident: restartTimes must survive the reset, so the breaker + // trips immediately at the restart threshold (3 restarts already in + // the rolling hour). + for i := 0; i < 5; i++ { + f.tick(channelhealth.OverallUnhealthy) + } + if len(f.restarts) != 3 { + t.Fatalf("restart log lost across incident reset: got %d restarts", len(f.restarts)) + } + // Second recovery fires exactly one more recovered event. + f.tick(channelhealth.OverallHealthy) + count := 0 + for _, tp := range f.eventTypes() { + if tp == eventRecovered { + count++ + } + } + if count != 2 { + t.Fatalf("expected 2 recovered events, got %d", count) + } +} + +func TestEscalator_RecoveryWithoutAlertIsSilent(t *testing.T) { + f := newEscalatorFixture(t) + f.tick(channelhealth.OverallUnhealthy) + f.tick(channelhealth.OverallHealthy) + if len(f.events) != 0 { + t.Fatalf("expected no events for sub-threshold blip, got %v", f.eventTypes()) + } +} + +func TestEscalator_DegradedHoldsIncident(t *testing.T) { + f := newEscalatorFixture(t) + f.tick(channelhealth.OverallUnhealthy) + f.tick(channelhealth.OverallUnhealthy) + // Degraded neither counts nor resets. + f.tick(channelhealth.OverallDegraded) + f.tick(channelhealth.OverallUnknown) + f.tick(channelhealth.OverallUnhealthy) + if got := f.eventTypes(); len(got) != 1 || got[0] != eventFailureDetected { + t.Fatalf("expected alert on 3rd failing check across hold, got %v", got) + } +} + +func TestEscalator_UnreachableCountsAsFailing(t *testing.T) { + f := newEscalatorFixture(t) + for i := 0; i < 3; i++ { + f.tick(channelhealth.OverallUnreachable) + } + if got := f.eventTypes(); len(got) != 1 || got[0] != eventFailureDetected { + t.Fatalf("expected alert for unreachable, got %v", got) + } +} + +func TestEscalator_ProductionRestartSkipsNonRunning(t *testing.T) { + f := newEscalatorFixture(t) + // Use the real restart dependency against a stopped instance row: it + // must no-op via restartInstanceAsyncWithToast's status guard. + inst := database.Instance{Name: "bot-x", DisplayName: "x", Status: "stopped"} + if err := database.DB.Create(&inst).Error; err != nil { + t.Fatalf("create instance: %v", err) + } + f.esc.restartInstance(inst.ID, "t", "m") + var got database.Instance + if err := database.DB.First(&got, inst.ID).Error; err != nil { + t.Fatalf("reload: %v", err) + } + if got.Status != "stopped" { + t.Fatalf("expected stopped instance untouched, got status %q", got.Status) + } +} diff --git a/control-plane/internal/handlers/channels.go b/control-plane/internal/handlers/channels.go index fde9e7e0..9494ad5d 100644 --- a/control-plane/internal/handlers/channels.go +++ b/control-plane/internal/handlers/channels.go @@ -125,3 +125,41 @@ func rfc3339OrNil(t *time.Time) *string { s := t.UTC().Format(time.RFC3339) return &s } + +// GetChannelHealthEvents returns the escalation audit log for an instance +// (alerts sent, auto-restarts, recoveries), newest first. +// GET /api/v1/instances/{id}/channels/health/events?limit=50 +func GetChannelHealthEvents(w http.ResponseWriter, r *http.Request) { + id, err := strconv.Atoi(chi.URLParam(r, "id")) + if err != nil { + writeError(w, http.StatusBadRequest, "Invalid instance ID") + return + } + + var inst database.Instance + if err := database.DB.First(&inst, id).Error; err != nil { + writeError(w, http.StatusNotFound, "Instance not found") + return + } + + if !middleware.CanAccessInstance(r, inst.ID) { + writeError(w, http.StatusForbidden, "Access denied") + return + } + + limit := 50 + if v := r.URL.Query().Get("limit"); v != "" { + if n, err := strconv.Atoi(v); err == nil && n > 0 && n <= 500 { + limit = n + } + } + + var events []database.ChannelHealthEvent + if err := database.DB.Where("instance_id = ?", inst.ID). + Order("created_at DESC").Limit(limit).Find(&events).Error; err != nil { + writeError(w, http.StatusInternalServerError, "Failed to load events") + return + } + + writeJSON(w, http.StatusOK, map[string]any{"events": events}) +} diff --git a/control-plane/internal/handlers/settings.go b/control-plane/internal/handlers/settings.go index b1cf70ce..427bfae5 100644 --- a/control-plane/internal/handlers/settings.go +++ b/control-plane/internal/handlers/settings.go @@ -12,7 +12,8 @@ import ( // fixedEncryptedSettings are non-LLM keys stored as fixed setting entries. var fixedEncryptedSettings = map[string]bool{ - "brave_api_key": true, + "brave_api_key": true, + "channel_alert_webhook_token": true, } // plainSettings are returned as-is (not encrypted). @@ -35,6 +36,9 @@ var plainSettings = []string{ "default_user_agent", "default_models", "analytics_consent", + "channel_alerts_enabled", + "channel_auto_restart_enabled", + "channel_alert_webhook_url", } func getAllSettings() map[string]string { @@ -186,6 +190,22 @@ func UpdateSettings(w http.ResponseWriter, r *http.Request) { } } + // Handle channel_alert_webhook_token (fixed encrypted) + if v, ok := raw["channel_alert_webhook_token"]; ok { + if strVal, ok := v.(string); ok { + if strVal != "" { + encrypted, err := utils.Encrypt(strVal) + if err != nil { + writeError(w, http.StatusInternalServerError, "Failed to encrypt webhook token") + return + } + database.SetSetting("channel_alert_webhook_token", encrypted) + } else { + database.SetSetting("channel_alert_webhook_token", "") + } + } + } + // Handle env_vars_set / env_vars_unset (PATCH-style for the encrypted map). // envVarsChanged is true only when the resulting plaintext map actually // differs from what was stored — a no-op request (e.g. re-setting the same @@ -233,7 +253,8 @@ func UpdateSettings(w http.ResponseWriter, r *http.Request) { // Handle remaining plain settings for key, val := range raw { - if key == "default_models" || key == "brave_api_key" || key == "env_vars_set" || key == "env_vars_unset" { + if key == "default_models" || key == "brave_api_key" || key == "channel_alert_webhook_token" || + key == "env_vars_set" || key == "env_vars_unset" { continue } if key == "default_pod_annotations" || key == "default_node_selector" || key == "default_tolerations" || diff --git a/control-plane/main.go b/control-plane/main.go index 3d5d60ea..ab1383b9 100644 --- a/control-plane/main.go +++ b/control-plane/main.go @@ -302,6 +302,15 @@ func main() { if config.Cfg.ChannelHealthEnabled { chMon := channelhealth.New(tunnelMgr, config.Cfg.ChannelHealthInterval) handlers.ChannelHealthMon = chMon + // Escalation: webhook alerts + opt-in auto-restart on sustained + // channel failure. Must be registered before Start. + esc := handlers.NewChannelEscalator(handlers.ChannelEscalatorConfig{ + AlertThreshold: config.Cfg.ChannelHealthAlertThreshold, + RestartThreshold: config.Cfg.ChannelHealthRestartThreshold, + MaxRestartsPerHour: config.Cfg.ChannelHealthRestartMaxPerHour, + RestartCooldown: config.Cfg.ChannelHealthRestartCooldown, + }) + chMon.SetListener(esc.OnSnapshot) chMon.Start(ctx) log.Printf("Channel health monitor started (interval=%s)", config.Cfg.ChannelHealthInterval) } @@ -379,6 +388,7 @@ func main() { r.Get("/instances/{id}/ssh-test", handlers.SSHConnectionTest) r.Get("/instances/{id}/ssh-status", handlers.GetSSHStatus) r.Get("/instances/{id}/channels/health", handlers.GetChannelHealth) + r.Get("/instances/{id}/channels/health/events", handlers.GetChannelHealthEvents) r.Get("/instances/{id}/ssh-events", handlers.GetSSHEvents) r.Post("/instances/{id}/ssh-reconnect", handlers.SSHReconnect) r.Get("/instances/{id}/tunnels", handlers.GetTunnelStatus) @@ -492,6 +502,7 @@ func main() { // Settings r.Get("/settings", handlers.GetSettings) r.Put("/settings", handlers.UpdateSettings) + r.Post("/settings/channel-alerts/test", handlers.TestChannelAlertWebhook) r.Post("/settings/rotate-ssh-key", handlers.RotateSSHKey) r.Get("/audit-logs", handlers.GetAuditLogs) diff --git a/docs/channel-monitoring.md b/docs/channel-monitoring.md index 217c477f..98ee615e 100644 --- a/docs/channel-monitoring.md +++ b/docs/channel-monitoring.md @@ -65,8 +65,52 @@ Overall instance status is derived from the per-channel states: | Env var | Default | Meaning | |---|---|---| -| `CLAWORC_CHANNEL_HEALTH_ENABLED` | `true` | Enable the poller | +| `CLAWORC_CHANNEL_HEALTH_ENABLED` | `true` | Enable the poller (and with it the escalation pipeline) | | `CLAWORC_CHANNEL_HEALTH_INTERVAL` | `60s` | Time between checks | +| `CLAWORC_CHANNEL_HEALTH_ALERT_THRESHOLD` | `3` | Consecutive failing checks before an alert fires | +| `CLAWORC_CHANNEL_HEALTH_RESTART_THRESHOLD` | `5` | Consecutive failing checks before an auto-restart fires | +| `CLAWORC_CHANNEL_HEALTH_RESTART_MAX_PER_HOUR` | `3` | Auto-restart circuit breaker (per instance, rolling hour) | +| `CLAWORC_CHANNEL_HEALTH_RESTART_COOLDOWN` | `10m` | After a triggered restart, failing checks are ignored for this long | + +Runtime behavior (UI-editable, stored in the settings table): + +| Setting key | Default | Meaning | +|---|---|---| +| `channel_alerts_enabled` | `true` | Deliver webhook alerts (inert without a URL) | +| `channel_alert_webhook_url` | empty | Where alert JSON is POSTed | +| `channel_alert_webhook_token` | empty | Optional `Authorization: Bearer` token (encrypted at rest) | +| `channel_auto_restart_enabled` | `false` | Opt-in automatic instance restarts | + +## Escalation + +The monitor feeds every snapshot to an escalator (`internal/handlers/channel_escalation.go`) +that tracks consecutive failing checks per instance. "Failing" means overall +`unhealthy` or `unreachable`; `degraded`/`unknown` *hold* an open incident +(neither count nor reset it); `healthy`/`no_channels` close it. + +Escalation ladder: + +1. **Alert** — at the alert threshold, one `channel_failure` webhook fires + per incident. +2. **Auto-restart** (opt-in) — at the restart threshold the instance is + restarted through the same async flow as a manual restart (tunnels + stopped, task + toast emitted). Guarded by the per-hour circuit breaker + and the post-restart cooldown; when the breaker trips, a single + `restart_limit_reached` webhook asks for manual intervention. +3. **Recovery** — when the incident closes after an alert was sent, one + `recovery` webhook reports the outage duration. + +Alert payloads are JSON with a human-readable `text` field plus structured +fields (`event`, `instance`, `overall`, `consecutive_failures`, +`failing_since`, `channels[]`). Delivery is fire-and-forget with one retry +on network error or 5xx. Admins can verify delivery with +`POST /api/v1/settings/channel-alerts/test` (the "Send Test" button in +Settings → Misc). + +Every escalation action is recorded in the `channel_health_events` audit +table and readable via `GET /api/v1/instances/{id}/channels/health/events`. +Incident counters are in-memory: a control-plane restart re-counts an +ongoing outage from zero (worst case, a duplicate alert after ~3 checks). ## API @@ -118,6 +162,17 @@ pair (one row per channel, overwritten on each check): | `Error` | string | Last error reported by the gateway, if any | | `CheckedAt` | datetime | When this status was recorded | +`channel_health_events` is the append-only escalation audit log: + +| Field | Type | Description | +|---|---|---| +| `InstanceID` | uint | Instance the event belongs to | +| `Type` | string | `failure_detected`, `auto_restart`, `restart_limit_reached`, `recovered`, `webhook_test` | +| `Overall` | string | Overall status at the time of the event | +| `Detail` | text | JSON context (failing channels, consecutive count, outage duration) | +| `WebhookStatus` | string | `sent`, `failed`, or `skipped` | +| `CreatedAt` | datetime | When the event occurred | + ## UI - **Channel Health panel** on the Agent detail page (Settings tab): @@ -134,12 +189,10 @@ pair (one row per channel, overwritten on each check): - Staleness is inferred from event silence, so a genuinely quiet channel (nobody messaging the agent for 30+ minutes) can be reported `stale` even though it is fine. -- Monitoring only observes — it does not restart channels, restart - instances, or deliver alerts. +- Escalation counters are in-memory only; a control-plane restart resets + consecutive-failure counts and the restart circuit-breaker window. ## Future work -- Auto-restart escalation for unhealthy channels. -- Alert delivery (notify operators when an agent goes unhealthy). - Synthetic canary probes to distinguish quiet channels from stale ones. - Consuming the gateway's push `health` broadcast instead of polling.