Fix alert and notification telemetry signals

This commit is contained in:
courtmanr@gmail.com
2026-07-29 13:44:49 +01:00
parent 000986f125
commit 848e166f5d
27 changed files with 818 additions and 75 deletions
+14 -1
View File
@@ -90,12 +90,19 @@ Every field is listed below with the reason it exists. Nothing else is included
| Notifications enabled | `true`/`false` | See whether alert notification delivery is configured |
| AI actions enabled | `true`/`false` | See whether AI control tools are enabled without sending action history or command content |
| Active alerts | `4` | Understand how noisy or quiet installations are in aggregate |
| Alerts fired 30d | `18` | Count locally retained alert-history entries in the current 30-day window without sending alert text, resource IDs, or timestamps |
| Alerts fired 30d | `18` | Count unique locally retained alert occurrences in the current 30-day window without sending alert text, resource IDs, or timestamps |
| Alerts acknowledged 30d | `7` | Count acknowledgements in the current 30-day window without sending actors, reasons, alert IDs, or timestamps |
| Alerts resolved 30d | `12` | Count resolved alert records in the current 30-day window without sending resolution details, alert IDs, or resource IDs |
| Notification attempts 7d | `14` | Count delivery attempts, including retry attempts, in the locally retained seven-day queue window without sending recipients, endpoints, titles, or message content |
| Notification deliveries 7d | `11` | Count successfully delivered queue records in the local seven-day window without sending channel, recipient, endpoint, or content |
| Notification failures 7d (terminal in schema v3) | `3` | Count terminal failed or dead-lettered delivery outcomes in the local seven-day window without sending retry-attempt failures, error text, endpoint, recipient, or message content |
| Notification failures authentication 7d (schema v5) | `1` | Count terminal failures classified locally as credential or permission failures without sending raw errors, provider names, endpoints, recipients, or message content |
| Notification failures rate limited 7d (schema v5) | `0` | Count terminal failures classified locally as provider rate limiting without sending provider identity or response content |
| Notification failures connectivity 7d (schema v5) | `2` | Count terminal failures classified locally as DNS, timeout, or connection failures without sending addresses, hosts, or raw errors |
| Notification failures TLS 7d (schema v5) | `0` | Count terminal failures classified locally as certificate or TLS failures without sending certificates, hostnames, or raw errors |
| Notification failures configuration 7d (schema v5) | `0` | Count terminal failures classified locally as missing or invalid destination configuration without sending configuration values |
| Notification failures rejected 7d (schema v5) | `0` | Count terminal failures classified locally as destination request or payload rejection without sending response or payload content |
| Notification failures unknown 7d (schema v5) | `0` | Count terminal failures that do not match another fixed class without sending raw errors |
| Relay enabled | `true`/`false` | See whether remote-access features are being used |
| SSO enabled | `true`/`false` | See whether single-sign-on support is being used |
| Multi-tenant | `true`/`false` | See whether multi-tenant/runtime-org features are being used |
@@ -186,6 +193,12 @@ pre-dispatch refusal categories, and an independently verified
action-to-finding resolution count. These remain aggregate counters and do not
send action, finding, resource, evidence, actor, or command identity or content.
Telemetry schema v5 adds seven fixed notification terminal-failure counters so
fleet health can distinguish authentication, rate limiting, connectivity, TLS,
configuration, destination rejection, and unknown failures. Classification is
performed locally; raw error text and destination/provider identity are never
included in the payload.
#### Server-side handling and retention
- Telemetry pings are stored on the Pulse license server only for aggregate install/use analysis.
+14 -1
View File
@@ -90,12 +90,19 @@ Every field is listed below with the reason it exists. Nothing else is included
| Notifications enabled | `true`/`false` | See whether alert notification delivery is configured |
| AI actions enabled | `true`/`false` | See whether AI control tools are enabled without sending action history or command content |
| Active alerts | `4` | Understand how noisy or quiet installations are in aggregate |
| Alerts fired 30d | `18` | Count locally retained alert-history entries in the current 30-day window without sending alert text, resource IDs, or timestamps |
| Alerts fired 30d | `18` | Count unique locally retained alert occurrences in the current 30-day window without sending alert text, resource IDs, or timestamps |
| Alerts acknowledged 30d | `7` | Count acknowledgements in the current 30-day window without sending actors, reasons, alert IDs, or timestamps |
| Alerts resolved 30d | `12` | Count resolved alert records in the current 30-day window without sending resolution details, alert IDs, or resource IDs |
| Notification attempts 7d | `14` | Count delivery attempts, including retry attempts, in the locally retained seven-day queue window without sending recipients, endpoints, titles, or message content |
| Notification deliveries 7d | `11` | Count successfully delivered queue records in the local seven-day window without sending channel, recipient, endpoint, or content |
| Notification failures 7d (terminal in schema v3) | `3` | Count terminal failed or dead-lettered delivery outcomes in the local seven-day window without sending retry-attempt failures, error text, endpoint, recipient, or message content |
| Notification failures authentication 7d (schema v5) | `1` | Count terminal failures classified locally as credential or permission failures without sending raw errors, provider names, endpoints, recipients, or message content |
| Notification failures rate limited 7d (schema v5) | `0` | Count terminal failures classified locally as provider rate limiting without sending provider identity or response content |
| Notification failures connectivity 7d (schema v5) | `2` | Count terminal failures classified locally as DNS, timeout, or connection failures without sending addresses, hosts, or raw errors |
| Notification failures TLS 7d (schema v5) | `0` | Count terminal failures classified locally as certificate or TLS failures without sending certificates, hostnames, or raw errors |
| Notification failures configuration 7d (schema v5) | `0` | Count terminal failures classified locally as missing or invalid destination configuration without sending configuration values |
| Notification failures rejected 7d (schema v5) | `0` | Count terminal failures classified locally as destination request or payload rejection without sending response or payload content |
| Notification failures unknown 7d (schema v5) | `0` | Count terminal failures that do not match another fixed class without sending raw errors |
| Relay enabled | `true`/`false` | See whether remote-access features are being used |
| SSO enabled | `true`/`false` | See whether single-sign-on support is being used |
| Multi-tenant | `true`/`false` | See whether multi-tenant/runtime-org features are being used |
@@ -186,6 +193,12 @@ pre-dispatch refusal categories, and an independently verified
action-to-finding resolution count. These remain aggregate counters and do not
send action, finding, resource, evidence, actor, or command identity or content.
Telemetry schema v5 adds seven fixed notification terminal-failure counters so
fleet health can distinguish authentication, rate limiting, connectivity, TLS,
configuration, destination rejection, and unknown failures. Classification is
performed locally; raw error text and destination/provider identity are never
included in the payload.
#### Server-side handling and retention
- Telemetry pings are stored on the Pulse license server only for aggregate install/use analysis.
@@ -106,6 +106,17 @@ describe('NotificationsAPI', () => {
counts_are_retention_bounded: true,
retry_attempts_affect_health: false,
terminal_failures_affect_health: true,
failure_classes_7d: {
authentication: 3,
rate_limited: 0,
connectivity: 1,
tls: 0,
configuration: 0,
rejected: 0,
unknown: 0,
},
failure_classes_available: true,
failure_class_window_days: 7,
},
} as any);
@@ -129,6 +140,17 @@ describe('NotificationsAPI', () => {
countsAreRetentionBounded: true,
retryAttemptsAffectHealth: false,
terminalFailuresAffectHealth: true,
failureClasses7d: {
authentication: 3,
rate_limited: 0,
connectivity: 1,
tls: 0,
configuration: 0,
rejected: 0,
unknown: 0,
},
failureClassesAvailable: true,
failureClassWindowDays: 7,
},
});
});
@@ -50,6 +50,8 @@ const mockTelemetryPreviewPayload = {
vmware_vms: 0,
vmware_datastores: 0,
availability_targets: 1,
availability_probe_targets: 0,
availability_probe_agents: 0,
ai_enabled: false,
patrol_enabled: false,
discovery_enabled: false,
@@ -71,6 +73,13 @@ const mockTelemetryPreviewPayload = {
notification_attempts_7d: 0,
notification_deliveries_7d: 0,
notification_failures_7d: 0,
notification_failures_authentication_7d: 0,
notification_failures_rate_limited_7d: 0,
notification_failures_connectivity_7d: 0,
notification_failures_tls_7d: 0,
notification_failures_configuration_7d: 0,
notification_failures_rejected_7d: 0,
notification_failures_unknown_7d: 0,
pulse_intelligence_loop_configured: false,
pulse_intelligence_loop_active_30d: false,
pulse_intelligence_complete_operations_loop_30d: false,
+53
View File
@@ -89,6 +89,16 @@ export interface NotificationTestRequest {
}
export type NotificationQueueHealthStatus = 'healthy' | 'degraded' | 'unavailable';
export type NotificationFailureClass =
| 'authentication'
| 'rate_limited'
| 'connectivity'
| 'tls'
| 'configuration'
| 'rejected'
| 'unknown';
export type NotificationFailureClassCounts = Record<NotificationFailureClass, number>;
export interface NotificationQueueHealth {
pending: number;
@@ -105,6 +115,9 @@ export interface NotificationQueueHealth {
countsAreRetentionBounded: boolean;
retryAttemptsAffectHealth: boolean;
terminalFailuresAffectHealth: boolean;
failureClasses7d: NotificationFailureClassCounts;
failureClassesAvailable: boolean;
failureClassWindowDays: number;
}
export interface NotificationHealth {
@@ -127,6 +140,27 @@ function nonNegativeCount(value: unknown): number | undefined {
: undefined;
}
const notificationFailureClasses: NotificationFailureClass[] = [
'authentication',
'rate_limited',
'connectivity',
'tls',
'configuration',
'rejected',
'unknown',
];
function normalizeFailureClassCounts(value: unknown): NotificationFailureClassCounts | undefined {
const record = apiRecord(value);
const counts = {} as NotificationFailureClassCounts;
for (const failureClass of notificationFailureClasses) {
const count = nonNegativeCount(record[failureClass]);
if (count === undefined) return undefined;
counts[failureClass] = count;
}
return counts;
}
export class NotificationsAPI {
private static baseUrl = '/api/notifications';
@@ -145,6 +179,9 @@ export class NotificationsAPI {
const attentionRequired = nonNegativeCount(queue.attention_required);
const completedRetentionDays = nonNegativeCount(queue.completed_retention_days);
const deadLetterRetentionDays = nonNegativeCount(queue.dead_letter_retention_days);
const failureClasses7d = normalizeFailureClassCounts(queue.failure_classes_7d);
const failureClassWindowDays = nonNegativeCount(queue.failure_class_window_days);
const failureClassesAvailable = strictBoolean(queue.failure_classes_available);
const reasonCodesAreValid =
Array.isArray(queue.reason_codes) &&
queue.reason_codes.every((reasonCode) => typeof reasonCode === 'string');
@@ -161,6 +198,9 @@ export class NotificationsAPI {
completedRetentionDays > 0 &&
deadLetterRetentionDays !== undefined &&
deadLetterRetentionDays > 0 &&
failureClasses7d !== undefined &&
failureClassWindowDays !== undefined &&
failureClassWindowDays > 0 &&
reasonCodesAreValid &&
semanticsAreValid;
const rawStatus = normalizeQueueHealthStatus(queue.status);
@@ -199,6 +239,19 @@ export class NotificationsAPI {
countsAreRetentionBounded: strictBoolean(queue.counts_are_retention_bounded),
retryAttemptsAffectHealth: strictBoolean(queue.retry_attempts_affect_health),
terminalFailuresAffectHealth: strictBoolean(queue.terminal_failures_affect_health, true),
failureClasses7d:
failureClasses7d ??
({
authentication: 0,
rate_limited: 0,
connectivity: 0,
tls: 0,
configuration: 0,
rejected: 0,
unknown: 0,
} satisfies NotificationFailureClassCounts),
failureClassesAvailable,
failureClassWindowDays: failureClassWindowDays ?? 7,
},
};
}
+9
View File
@@ -51,6 +51,8 @@ export interface TelemetryPingPreview {
vmware_vms: number;
vmware_datastores: number;
availability_targets: number;
availability_probe_targets: number;
availability_probe_agents: number;
ai_enabled: boolean;
patrol_enabled: boolean;
discovery_enabled: boolean;
@@ -72,6 +74,13 @@ export interface TelemetryPingPreview {
notification_attempts_7d: number;
notification_deliveries_7d: number;
notification_failures_7d: number;
notification_failures_authentication_7d: number;
notification_failures_rate_limited_7d: number;
notification_failures_connectivity_7d: number;
notification_failures_tls_7d: number;
notification_failures_configuration_7d: number;
notification_failures_rejected_7d: number;
notification_failures_unknown_7d: number;
pulse_intelligence_loop_configured: boolean;
pulse_intelligence_loop_active_30d: boolean;
pulse_intelligence_complete_operations_loop_30d: boolean;
@@ -56,6 +56,8 @@ const buildTelemetryPreviewPayload = (
vmware_vms: 0,
vmware_datastores: 0,
availability_targets: 1,
availability_probe_targets: 0,
availability_probe_agents: 0,
ai_enabled: false,
patrol_enabled: false,
discovery_enabled: false,
@@ -77,6 +79,13 @@ const buildTelemetryPreviewPayload = (
notification_attempts_7d: 0,
notification_deliveries_7d: 0,
notification_failures_7d: 0,
notification_failures_authentication_7d: 0,
notification_failures_rate_limited_7d: 0,
notification_failures_connectivity_7d: 0,
notification_failures_tls_7d: 0,
notification_failures_configuration_7d: 0,
notification_failures_rejected_7d: 0,
notification_failures_unknown_7d: 0,
pulse_intelligence_loop_configured: false,
pulse_intelligence_loop_active_30d: false,
pulse_intelligence_complete_operations_loop_30d: false,
@@ -20,6 +20,17 @@ const degradedHealth: NotificationQueueHealth = {
countsAreRetentionBounded: true,
retryAttemptsAffectHealth: false,
terminalFailuresAffectHealth: true,
failureClasses7d: {
authentication: 2,
rate_limited: 0,
connectivity: 1,
tls: 0,
configuration: 0,
rejected: 0,
unknown: 0,
},
failureClassesAvailable: true,
failureClassWindowDays: 7,
};
describe('AlertDeliveryHealthCard', () => {
@@ -42,7 +53,13 @@ describe('AlertDeliveryHealthCard', () => {
'2 dead-lettered deliveries retained for 30 days',
);
expect(screen.getByRole('alert')).toHaveTextContent(
'recoverable retry attempts do not trigger this warning',
'Recoverable retry attempts do not trigger this warning',
);
expect(screen.getByRole('alert')).toHaveTextContent(
'classified as authentication (2)',
);
expect(screen.getByRole('alert')).toHaveTextContent(
'Check destination credentials, tokens, and account permissions',
);
fireEvent.click(screen.getByRole('button', { name: 'Refresh delivery status' }));
@@ -40,6 +40,8 @@ export function AlertDeliveryHealthCard(props: AlertDeliveryHealthCardProps) {
deadLetter: deadLetter(),
completedRetentionDays: completedRetentionDays(),
deadLetterRetentionDays: deadLetterRetentionDays(),
failureClasses7d: props.health?.failureClasses7d,
failureClassesAvailable: props.health?.failureClassesAvailable ?? false,
})}
</p>
</div>
@@ -117,6 +117,17 @@ describe('useAlertDestinationsTabState', () => {
countsAreRetentionBounded: true,
retryAttemptsAffectHealth: false,
terminalFailuresAffectHealth: true,
failureClasses7d: {
authentication: 0,
rate_limited: 0,
connectivity: 0,
tls: 0,
configuration: 0,
rejected: 0,
unknown: 0,
},
failureClassesAvailable: true,
failureClassWindowDays: 7,
},
});
vi.mocked(NotificationsAPI.testNotification).mockResolvedValue({ success: true } as never);
@@ -133,9 +133,19 @@ describe('alertDestinationsPresentation', () => {
deadLetter: 2,
completedRetentionDays: 7,
deadLetterRetentionDays: 30,
failureClasses7d: {
authentication: 3,
rate_limited: 0,
connectivity: 0,
tls: 0,
configuration: 0,
rejected: 0,
unknown: 0,
},
failureClassesAvailable: true,
}),
).toBe(
'1 failed delivery retained for 7 days and 2 dead-lettered deliveries retained for 30 days. These notifications were not delivered. Check each enabled destination and send a test; recoverable retry attempts do not trigger this warning.',
'1 failed delivery retained for 7 days and 2 dead-lettered deliveries retained for 30 days. These notifications were not delivered. Most recent terminal failures were classified as authentication (3). Check destination credentials, tokens, and account permissions. Recoverable retry attempts do not trigger this warning.',
);
expect(
getAlertDestinationsDeliveryHealthDescription({
@@ -154,6 +154,8 @@ export function getAlertDestinationsDeliveryHealthDescription(input: {
deadLetter: number;
completedRetentionDays: number;
deadLetterRetentionDays: number;
failureClasses7d?: Record<string, number>;
failureClassesAvailable?: boolean;
}) {
if (input.status === 'unavailable') {
return 'Pulse could not verify the notification queue. Review the destination settings below and send a test before relying on delivery.';
@@ -171,7 +173,25 @@ export function getAlertDestinationsDeliveryHealthDescription(input: {
);
}
const summary = outcomes.join(' and ') || 'A retained terminal delivery failure';
return `${summary}. These notifications were not delivered. Check each enabled destination and send a test; recoverable retry attempts do not trigger this warning.`;
const guidanceByClass: Record<string, string> = {
authentication: 'Check destination credentials, tokens, and account permissions.',
rate_limited: 'Check provider rate limits and reduce delivery volume before retrying.',
connectivity: 'Check DNS, firewall, proxy, and destination reachability.',
tls: 'Check certificate trust, hostname matching, and TLS settings.',
configuration: 'Review the enabled destination configuration and required fields.',
rejected: 'Check the destination endpoint and payload requirements.',
unknown: 'Review the local notification audit details for the terminal error.',
};
let diagnostic = 'Check each enabled destination and send a test.';
if (input.failureClassesAvailable && input.failureClasses7d) {
const dominant = Object.entries(input.failureClasses7d)
.filter(([, count]) => count > 0)
.sort((left, right) => right[1] - left[1])[0];
if (dominant) {
diagnostic = `Most recent terminal failures were classified as ${dominant[0].replace('_', ' ')} (${dominant[1]}). ${guidanceByClass[dominant[0]] ?? guidanceByClass.unknown}`;
}
}
return `${summary}. These notifications were not delivered. ${diagnostic} Recoverable retry attempts do not trigger this warning.`;
}
export function getAlertDestinationsDeliveryRefreshLabel() {
+8 -5
View File
@@ -60,27 +60,30 @@ func (m *Manager) pruneRecentlyResolvedUnlocked(now time.Time) {
}
// consumeRecentlyResolvedForRefireWithPrimaryLock consumes a resolved alert
// that is still inside the refire window. The caller must hold m.mu.
// that is still inside the refire window. The final return value reports that
// a prior resolved occurrence exists even when it is too old to reactivate, so
// callers can explicitly append the genuine new occurrence to history.
// The caller must hold m.mu.
//
// Lock order is always m.mu -> resolvedMutex. This helper performs only
// resolved-map access while resolvedMutex is held; history, dispatch, and
// notification work remain the caller's responsibility.
func (m *Manager) consumeRecentlyResolvedForRefireWithPrimaryLock(storageKey string, now time.Time) (time.Time, time.Time, bool) {
func (m *Manager) consumeRecentlyResolvedForRefireWithPrimaryLock(storageKey string, now time.Time) (time.Time, time.Time, bool, bool) {
m.resolvedMutex.Lock()
defer m.resolvedMutex.Unlock()
resolved, ok := m.getResolvedAlertNoLock(storageKey)
if !ok || resolved == nil || resolved.Alert == nil {
return time.Time{}, time.Time{}, false
return time.Time{}, time.Time{}, false, false
}
if !resolved.ResolvedTime.After(now.Add(-recentlyResolvedRetention)) {
return time.Time{}, time.Time{}, false
return time.Time{}, resolved.ResolvedTime, false, true
}
startTime := resolved.Alert.StartTime
resolvedAt := resolved.ResolvedTime
m.removeResolvedAlertUnlocked(storageKey)
return startTime, resolvedAt, true
return startTime, resolvedAt, true, true
}
// addRecentlyResolvedWithPrimaryLock records a resolved alert while preserving the caller's
+13 -5
View File
@@ -445,7 +445,7 @@ func (m *Manager) evaluateCanonicalLifecycleAlert(params canonicalLifecycleAlert
return result, true
}
reactivatedStart, reactivatedAt, reactivated := m.consumeRecentlyResolvedForRefireWithPrimaryLock(storageKey, time.Now())
reactivatedStart, reactivatedAt, reactivated, hadResolvedOccurrence := m.consumeRecentlyResolvedForRefireWithPrimaryLock(storageKey, time.Now())
if reactivated {
if !reactivatedStart.IsZero() {
alert.StartTime = reactivatedStart
@@ -465,7 +465,11 @@ func (m *Manager) evaluateCanonicalLifecycleAlert(params canonicalLifecycleAlert
}
if params.AddToHistory {
m.historyManager.AddAlert(*alert)
if hadResolvedOccurrence {
m.historyManager.AddAlertTransition(*alert)
} else {
m.historyManager.AddAlert(*alert)
}
}
if params.RateLimit && !m.checkRateLimit(trackingKey) {
@@ -606,7 +610,7 @@ func (m *Manager) evaluateCanonicalStatefulAlert(params canonicalStatefulAlertPa
}
if existing == nil {
reactivatedStart, reactivatedAt, reactivated := m.consumeRecentlyResolvedForRefireWithPrimaryLock(storageKey, time.Now())
reactivatedStart, reactivatedAt, reactivated, hadResolvedOccurrence := m.consumeRecentlyResolvedForRefireWithPrimaryLock(storageKey, time.Now())
if reactivated {
if !reactivatedStart.IsZero() {
alert.StartTime = reactivatedStart
@@ -626,7 +630,11 @@ func (m *Manager) evaluateCanonicalStatefulAlert(params canonicalStatefulAlertPa
}
if params.AddToHistory {
m.historyManager.AddAlert(*alert)
if hadResolvedOccurrence {
m.historyManager.AddAlertTransition(*alert)
} else {
m.historyManager.AddAlert(*alert)
}
}
if params.RateLimit && !m.checkRateLimit(trackingKey) {
log.Debug().
@@ -642,7 +650,7 @@ func (m *Manager) evaluateCanonicalStatefulAlert(params canonicalStatefulAlertPa
if result.Transition != nil && result.Transition.Kind == alertspecs.EvaluationTransitionSeverityChanged && params.NotifyOnSeverityChange {
if params.AddToHistoryOnSeverityChange {
m.historyManager.AddAlert(*alert)
m.historyManager.AddAlertTransition(*alert)
}
if params.RateLimit && !m.checkRateLimit(trackingKey) {
log.Debug().
+82 -4
View File
@@ -9,6 +9,7 @@ import (
"sync"
"time"
"github.com/rcourtman/pulse-go-rewrite/internal/operationaltrust"
"github.com/rcourtman/pulse-go-rewrite/internal/securityutil"
"github.com/rcourtman/pulse-go-rewrite/internal/utils"
"github.com/rs/zerolog/log"
@@ -135,8 +136,22 @@ func (hm *HistoryManager) OnAlert(cb AlertCallback) {
hm.callbacks = append(hm.callbacks, cb)
}
// AddAlert adds an alert to history
// AddAlert records a newly observed alert occurrence. Repeated attempts to add
// the same still-open occurrence are folded into its existing history row.
// This makes the history write path enforce the same incident semantics as the
// runtime alert lifecycle instead of relying on the daily cleanup pass.
func (hm *HistoryManager) AddAlert(alert Alert) {
hm.addAlert(alert, false)
}
// AddAlertTransition records an explicit transition snapshot even when the
// alert occurrence is still open. Callers should use this only when the extra
// row is intentional, such as a configured severity-change history event.
func (hm *HistoryManager) AddAlertTransition(alert Alert) {
hm.addAlert(alert, true)
}
func (hm *HistoryManager) addAlert(alert Alert, forceAppend bool) {
hm.mu.Lock()
entry := HistoryEntry{
@@ -144,6 +159,21 @@ func (hm *HistoryManager) AddAlert(alert Alert) {
Timestamp: time.Now(),
}
if !forceAppend {
for i := len(hm.history) - 1; i >= 0; i-- {
if historyIdentityKey(&hm.history[i].Alert) != historyIdentityKey(&entry.Alert) {
continue
}
if sameHistoryIncident(hm.history[i], entry) {
hm.history[i].Alert = mergeHistoryAlertSnapshots(hm.history[i].Alert, entry.Alert)
hm.mu.Unlock()
log.Debug().Str("alertID", alert.ID).Msg("updated existing alert history occurrence")
return
}
break
}
}
hm.history = append(hm.history, entry)
callbacks := append([]AlertCallback(nil), hm.callbacks...)
hm.mu.Unlock()
@@ -156,6 +186,56 @@ func (hm *HistoryManager) AddAlert(alert Alert) {
}
}
func sameHistoryIncident(existing, incoming HistoryEntry) bool {
existingKey := historyIdentityKey(&existing.Alert)
if existingKey == "" || existingKey != historyIdentityKey(&incoming.Alert) {
return false
}
if !existing.Alert.StartTime.IsZero() &&
!incoming.Alert.StartTime.IsZero() &&
existing.Alert.StartTime.Equal(incoming.Alert.StartTime) {
return true
}
existingResolved := existing.Alert.OperationalRecord != nil &&
existing.Alert.OperationalRecord.State == operationaltrust.OperationalResolved
incomingResolved := incoming.Alert.OperationalRecord != nil &&
incoming.Alert.OperationalRecord.State == operationaltrust.OperationalResolved
if existingResolved || incomingResolved {
return false
}
if existing.Alert.OperationalRecord != nil &&
incoming.Alert.OperationalRecord != nil &&
!existingResolved && !incomingResolved {
return true
}
lastObservation := existing.Timestamp
if existing.Alert.LastSeen.After(lastObservation) {
lastObservation = existing.Alert.LastSeen
}
gap := incoming.Timestamp.Sub(lastObservation)
return gap >= 0 && gap <= historyDedupWindow
}
func mergeHistoryAlertSnapshots(existing, incoming Alert) Alert {
merged := *incoming.Clone()
if !existing.StartTime.IsZero() &&
(merged.StartTime.IsZero() || existing.StartTime.Before(merged.StartTime)) {
merged.StartTime = existing.StartTime
}
if existing.LastSeen.After(merged.LastSeen) {
merged.LastSeen = existing.LastSeen
}
if existing.Acknowledged && !merged.Acknowledged {
merged.Acknowledged = true
merged.AckTime = existing.AckTime
merged.AckUser = existing.AckUser
}
return merged
}
// UpdateAlertLastSeen updates the LastSeen timestamp on the most recent
// history entry matching the given alert ID. This is called when an alert is
// resolved so that the stored history reflects the true duration of the alert,
@@ -630,9 +710,7 @@ func (hm *HistoryManager) deduplicateHistory() {
if lastTime, ok := lastTimeByKey[key]; ok {
if entry.Timestamp.Sub(lastTime) <= historyDedupWindow {
idx := lastIdxByKey[key]
if entry.Alert.LastSeen.After(deduped[idx].Alert.LastSeen) {
deduped[idx].Alert.LastSeen = entry.Alert.LastSeen
}
deduped[idx].Alert = mergeHistoryAlertSnapshots(deduped[idx].Alert, entry.Alert)
lastTimeByKey[key] = entry.Timestamp
removed++
continue
+79
View File
@@ -6,6 +6,8 @@ import (
"strings"
"testing"
"time"
"github.com/rcourtman/pulse-go-rewrite/internal/operationaltrust"
)
// newTestHistoryManager creates a HistoryManager using a temp directory
@@ -159,6 +161,83 @@ func TestAddAlert(t *testing.T) {
}
}
func TestAddAlertCoalescesRepeatedOpenOccurrence(t *testing.T) {
hm := newTestHistoryManager(t)
startedAt := time.Now().Add(-time.Hour)
first := Alert{
ID: "resource::cpu",
CanonicalState: "resource::cpu",
StartTime: startedAt,
LastSeen: startedAt.Add(time.Minute),
OperationalRecord: &operationaltrust.OperationalRecord{
State: operationaltrust.OperationalOpen,
},
}
recreated := *first.Clone()
recreated.StartTime = startedAt.Add(30 * time.Minute)
recreated.LastSeen = startedAt.Add(31 * time.Minute)
callbacks := 0
hm.OnAlert(func(Alert) { callbacks++ })
hm.AddAlert(first)
hm.AddAlert(recreated)
if len(hm.history) != 1 {
t.Fatalf("history length = %d, want one open occurrence", len(hm.history))
}
if !hm.history[0].Alert.StartTime.Equal(startedAt) {
t.Fatalf("coalesced start = %v, want %v", hm.history[0].Alert.StartTime, startedAt)
}
if !hm.history[0].Alert.LastSeen.Equal(recreated.LastSeen) {
t.Fatalf("coalesced last seen = %v, want %v", hm.history[0].Alert.LastSeen, recreated.LastSeen)
}
if callbacks != 1 {
t.Fatalf("callbacks = %d, want one new-incident callback", callbacks)
}
}
func TestAddAlertAppendsAfterResolvedOccurrence(t *testing.T) {
hm := newTestHistoryManager(t)
startedAt := time.Now().Add(-time.Hour)
resolved := Alert{
ID: "resource::cpu",
CanonicalState: "resource::cpu",
StartTime: startedAt,
OperationalRecord: &operationaltrust.OperationalRecord{
State: operationaltrust.OperationalResolved,
},
}
refired := *resolved.Clone()
refired.StartTime = startedAt.Add(30 * time.Minute)
refired.OperationalRecord.State = operationaltrust.OperationalOpen
hm.AddAlert(resolved)
hm.AddAlert(refired)
if len(hm.history) != 2 {
t.Fatalf("history length = %d, want distinct resolved and re-fired occurrences", len(hm.history))
}
}
func TestAddAlertTransitionExplicitlyPreservesSeverityHistory(t *testing.T) {
hm := newTestHistoryManager(t)
alert := Alert{
ID: "resource::cpu",
CanonicalState: "resource::cpu",
StartTime: time.Now().Add(-time.Hour),
OperationalRecord: &operationaltrust.OperationalRecord{
State: operationaltrust.OperationalOpen,
},
}
hm.AddAlert(alert)
hm.AddAlertTransition(alert)
if len(hm.history) != 2 {
t.Fatalf("history length = %d, want explicit transition snapshot", len(hm.history))
}
}
func TestUpdateAlertLastSeenForAlertMatchesCanonicalState(t *testing.T) {
hm := newTestHistoryManager(t)
+19
View File
@@ -9,6 +9,7 @@ import (
"net/http"
"strings"
"sync"
"time"
"github.com/rcourtman/pulse-go-rewrite/internal/config"
"github.com/rcourtman/pulse-go-rewrite/internal/monitoring"
@@ -43,6 +44,7 @@ type NotificationManager interface {
GetWebhookHistory() []notifications.WebhookDelivery
TestEnhancedWebhook(notifications.EnhancedWebhookConfig) (int, string, error)
GetQueueStats() (map[string]int, error)
GetTelemetryStats(time.Time) (notifications.TelemetryStats, error)
}
// NotificationConfigPersistence defines the interface for saving notification configuration.
@@ -793,6 +795,17 @@ func (h *NotificationHandlers) GetNotificationHealth(w http.ResponseWriter, r *h
monitor := h.getMonitor(r.Context())
manager := monitor.GetNotificationManager()
failureClasses := notifications.NotificationFailureClassCounts{}.AsMap()
failureClassesAvailable := false
if deliveryStats, statsErr := manager.GetTelemetryStats(
time.Now().Add(-completedQueueRetentionDays * 24 * time.Hour),
); statsErr != nil {
log.Warn().Err(statsErr).Msg("Failed to get notification failure classes for health check")
} else {
failureClasses = deliveryStats.FailureClasses.AsMap()
failureClassesAvailable = true
}
queueStats := make(map[string]interface{})
stats, err := manager.GetQueueStats()
if err != nil {
@@ -807,6 +820,9 @@ func (h *NotificationHandlers) GetNotificationHealth(w http.ResponseWriter, r *h
"counts_are_retention_bounded": true,
"retry_attempts_affect_health": false,
"terminal_failures_affect_health": true,
"failure_classes_7d": failureClasses,
"failure_classes_available": failureClassesAvailable,
"failure_class_window_days": completedQueueRetentionDays,
}
} else {
healthy, attentionRequired, reasonCodes := classifyNotificationQueueHealth(stats)
@@ -829,6 +845,9 @@ func (h *NotificationHandlers) GetNotificationHealth(w http.ResponseWriter, r *h
"counts_are_retention_bounded": true,
"retry_attempts_affect_health": false,
"terminal_failures_affect_health": true,
"failure_classes_7d": failureClasses,
"failure_classes_available": failureClassesAvailable,
"failure_class_window_days": completedQueueRetentionDays,
}
}
+30 -9
View File
@@ -9,6 +9,7 @@ import (
"testing"
"github.com/rcourtman/pulse-go-rewrite/internal/notifications"
"github.com/stretchr/testify/mock"
)
func TestGetNotificationHealthReportsRetainedTerminalFailures(t *testing.T) {
@@ -24,6 +25,13 @@ func TestGetNotificationHealthReportsRetainedTerminalFailures(t *testing.T) {
"failed": 2,
"dlq": 3,
}, nil).Once()
mockManager.On("GetTelemetryStats", mock.Anything).Return(notifications.TelemetryStats{
Failures: 5,
FailureClasses: notifications.NotificationFailureClassCounts{
Authentication: 3,
Connectivity: 2,
},
}, nil).Once()
mockManager.On("GetEmailConfig").Return(notifications.EmailConfig{}).Once()
mockManager.On("GetWebhooks").Return([]notifications.WebhookConfig{}).Once()
mockPersistence.On("IsEncryptionEnabled").Return(true).Once()
@@ -40,15 +48,18 @@ func TestGetNotificationHealthReportsRetainedTerminalFailures(t *testing.T) {
var response struct {
OverallHealthy bool `json:"overall_healthy"`
Queue struct {
Healthy bool `json:"healthy"`
Status string `json:"status"`
AttentionRequired int `json:"attention_required"`
ReasonCodes []string `json:"reason_codes"`
CompletedRetentionDays int `json:"completed_retention_days"`
DeadLetterRetentionDays int `json:"dead_letter_retention_days"`
CountsAreRetentionBounded bool `json:"counts_are_retention_bounded"`
RetryAttemptsAffectHealth bool `json:"retry_attempts_affect_health"`
TerminalFailuresAffectHealth bool `json:"terminal_failures_affect_health"`
Healthy bool `json:"healthy"`
Status string `json:"status"`
AttentionRequired int `json:"attention_required"`
ReasonCodes []string `json:"reason_codes"`
CompletedRetentionDays int `json:"completed_retention_days"`
DeadLetterRetentionDays int `json:"dead_letter_retention_days"`
CountsAreRetentionBounded bool `json:"counts_are_retention_bounded"`
RetryAttemptsAffectHealth bool `json:"retry_attempts_affect_health"`
TerminalFailuresAffectHealth bool `json:"terminal_failures_affect_health"`
FailureClasses7d map[string]int `json:"failure_classes_7d"`
FailureClassesAvailable bool `json:"failure_classes_available"`
FailureClassWindowDays int `json:"failure_class_window_days"`
} `json:"queue"`
}
if err := json.Unmarshal(rec.Body.Bytes(), &response); err != nil {
@@ -72,6 +83,12 @@ func TestGetNotificationHealthReportsRetainedTerminalFailures(t *testing.T) {
!response.Queue.TerminalFailuresAffectHealth {
t.Fatalf("queue semantics = %#v", response.Queue)
}
if !response.Queue.FailureClassesAvailable ||
response.Queue.FailureClassWindowDays != 7 ||
response.Queue.FailureClasses7d["authentication"] != 3 ||
response.Queue.FailureClasses7d["connectivity"] != 2 {
t.Fatalf("failure classes = %#v", response.Queue)
}
}
func TestGetNotificationHealthFailsClosedWhenQueueStatsAreUnavailable(t *testing.T) {
@@ -81,6 +98,10 @@ func TestGetNotificationHealthFailsClosedWhenQueueStatsAreUnavailable(t *testing
mockMonitor.On("GetNotificationManager").Return(mockManager)
mockMonitor.On("GetConfigPersistence").Return(mockPersistence)
mockManager.On("GetQueueStats").Return(nil, errors.New("database path /secret unavailable")).Once()
mockManager.On("GetTelemetryStats", mock.Anything).Return(
notifications.TelemetryStats{},
errors.New("database unavailable"),
).Once()
mockManager.On("GetEmailConfig").Return(notifications.EmailConfig{}).Once()
mockManager.On("GetWebhooks").Return([]notifications.WebhookConfig{}).Once()
mockPersistence.On("IsEncryptionEnabled").Return(false).Once()
+14
View File
@@ -10,6 +10,7 @@ import (
"strings"
"sync"
"testing"
"time"
"github.com/rcourtman/pulse-go-rewrite/internal/notifications"
"github.com/stretchr/testify/assert"
@@ -252,6 +253,11 @@ func (m *MockNotificationManager) GetQueueStats() (map[string]int, error) {
return args.Get(0).(map[string]int), args.Error(1)
}
func (m *MockNotificationManager) GetTelemetryStats(since time.Time) (notifications.TelemetryStats, error) {
args := m.Called(since)
return args.Get(0).(notifications.TelemetryStats), args.Error(1)
}
type MockNotificationConfigPersistence struct {
mock.Mock
}
@@ -413,6 +419,10 @@ func TestNotificationHandlers(t *testing.T) {
"dlq": 0,
}
mockManager.On("GetQueueStats").Return(stats, nil).Once()
mockManager.On("GetTelemetryStats", mock.Anything).Return(
notifications.TelemetryStats{},
nil,
).Once()
mockManager.On("GetEmailConfig").Return(notifications.EmailConfig{}).Once()
mockManager.On("GetWebhooks").Return([]notifications.WebhookConfig{}).Once()
mockPersistence.On("IsEncryptionEnabled").Return(true).Once()
@@ -567,6 +577,10 @@ func TestNotificationHandlers(t *testing.T) {
{"GET", "/api/notifications/email-providers", func() {}},
{"GET", "/api/notifications/health", func() {
mockManager.On("GetQueueStats").Return(map[string]int{}, nil).Once()
mockManager.On("GetTelemetryStats", mock.Anything).Return(
notifications.TelemetryStats{},
nil,
).Once()
mockManager.On("GetEmailConfig").Return(notifications.EmailConfig{}).Once()
mockManager.On("GetWebhooks").Return([]notifications.WebhookConfig{}).Once()
mockPersistence.On("IsEncryptionEnabled").Return(true).Once()
+74 -33
View File
@@ -18,36 +18,43 @@ import (
// InstallSnapshotCounts holds install-wide resource and alert counts aggregated
// across tenant monitors.
type InstallSnapshotCounts struct {
PVENodes int
PBSInstances int
PMGInstances int
VMs int
Containers int
AgentHosts int
DockerHosts int
DockerContainers int
KubernetesClusters int
KubernetesNodes int
KubernetesPods int
KubernetesDeployments int
StoragePools int
PhysicalDisks int
CephClusters int
NetworkShares int
TrueNASSystems int
TrueNASVMs int
TrueNASApps int
VMwareHosts int
VMwareVMs int
VMwareDatastores int
AvailabilityTargets int
ActiveAlerts int
AlertsFired30d int
AlertsAcknowledged30d int
AlertsResolved30d int
NotificationAttempts7d int
NotificationDeliveries7d int
NotificationFailures7d int
PVENodes int
PBSInstances int
PMGInstances int
VMs int
Containers int
AgentHosts int
DockerHosts int
DockerContainers int
KubernetesClusters int
KubernetesNodes int
KubernetesPods int
KubernetesDeployments int
StoragePools int
PhysicalDisks int
CephClusters int
NetworkShares int
TrueNASSystems int
TrueNASVMs int
TrueNASApps int
VMwareHosts int
VMwareVMs int
VMwareDatastores int
AvailabilityTargets int
ActiveAlerts int
AlertsFired30d int
AlertsAcknowledged30d int
AlertsResolved30d int
NotificationAttempts7d int
NotificationDeliveries7d int
NotificationFailures7d int
NotificationFailuresAuthentication7d int
NotificationFailuresRateLimited7d int
NotificationFailuresConnectivity7d int
NotificationFailuresTLS7d int
NotificationFailuresConfiguration7d int
NotificationFailuresRejected7d int
NotificationFailuresUnknown7d int
}
// ReloadableMonitor wraps a Monitor with reload capability
@@ -289,6 +296,13 @@ func accumulateInstallOutcomeCounts(counts *InstallSnapshotCounts, monitor *Moni
counts.NotificationAttempts7d += stats.Attempts
counts.NotificationDeliveries7d += stats.Deliveries
counts.NotificationFailures7d += stats.Failures
counts.NotificationFailuresAuthentication7d += stats.FailureClasses.Authentication
counts.NotificationFailuresRateLimited7d += stats.FailureClasses.RateLimited
counts.NotificationFailuresConnectivity7d += stats.FailureClasses.Connectivity
counts.NotificationFailuresTLS7d += stats.FailureClasses.TLS
counts.NotificationFailuresConfiguration7d += stats.FailureClasses.Configuration
counts.NotificationFailuresRejected7d += stats.FailureClasses.Rejected
counts.NotificationFailuresUnknown7d += stats.FailureClasses.Unknown
}
}
@@ -296,12 +310,39 @@ func accumulateAlertOutcomeCounts(counts *InstallSnapshotCounts, history []alert
if counts == nil {
return
}
for _, alert := range history {
counts.AlertsFired30d++
type outcome struct {
acknowledged bool
resolved bool
}
outcomes := make(map[string]outcome, len(history))
for index, alert := range history {
identity := strings.TrimSpace(alert.CanonicalState)
if identity == "" {
identity = strings.TrimSpace(alert.ID)
}
occurrence := alert.StartTime
if occurrence.IsZero() && alert.OperationalRecord != nil {
occurrence = alert.OperationalRecord.FirstObservedAt
}
key := fmt.Sprintf("%s:%d", identity, occurrence.UnixNano())
if identity == "" {
key = fmt.Sprintf("anonymous:%d", index)
}
current := outcomes[key]
if alert.AckTime != nil && !alert.AckTime.Before(cutoff) {
counts.AlertsAcknowledged30d++
current.acknowledged = true
}
if alert.OperationalRecord != nil && alert.OperationalRecord.State == operationaltrust.OperationalResolved {
current.resolved = true
}
outcomes[key] = current
}
for _, current := range outcomes {
counts.AlertsFired30d++
if current.acknowledged {
counts.AlertsAcknowledged30d++
}
if current.resolved {
counts.AlertsResolved30d++
}
}
+28
View File
@@ -188,6 +188,34 @@ func TestAccumulateAlertOutcomeCountsUsesOnlyContentFreeLifecycleTotals(t *testi
assert.Equal(t, 1, counts.AlertsResolved30d)
}
func TestAccumulateAlertOutcomeCountsDeduplicatesOccurrenceSnapshots(t *testing.T) {
now := time.Now().UTC()
startedAt := now.Add(-2 * time.Hour)
acknowledgedAt := now.Add(-time.Hour)
open := alerts.Alert{
ID: "resource::cpu",
CanonicalState: "resource::cpu",
StartTime: startedAt,
}
acknowledged := *open.Clone()
acknowledged.AckTime = &acknowledgedAt
resolved := *acknowledged.Clone()
resolved.OperationalRecord = &operationaltrust.OperationalRecord{
State: operationaltrust.OperationalResolved,
}
counts := InstallSnapshotCounts{}
accumulateAlertOutcomeCounts(
&counts,
[]alerts.Alert{open, acknowledged, resolved},
now.Add(-30*24*time.Hour),
)
assert.Equal(t, 1, counts.AlertsFired30d)
assert.Equal(t, 1, counts.AlertsAcknowledged30d)
assert.Equal(t, 1, counts.AlertsResolved30d)
}
func testTelemetryMonitor(
nodes []models.Node,
vms []models.VM,
+149 -5
View File
@@ -27,6 +27,7 @@ const (
notificationAuditAlertIdentifiersColumn = "alert_identifiers"
legacyNotificationAuditAlertIdentifiersColumn = "alert_ids"
notificationOperationalLinksColumn = "operational_links"
notificationFailureClassColumn = "failure_class"
notificationQueueDirName = "notifications"
notificationQueueFileName = "notification_queue.db"
)
@@ -43,6 +44,97 @@ const (
QueueStatusCancelled NotificationQueueStatus = "cancelled"
)
// NotificationFailureClass is a closed, content-free delivery failure bucket.
// Only these coarse values may leave the installation in telemetry.
type NotificationFailureClass string
const (
NotificationFailureAuthentication NotificationFailureClass = "authentication"
NotificationFailureRateLimited NotificationFailureClass = "rate_limited"
NotificationFailureConnectivity NotificationFailureClass = "connectivity"
NotificationFailureTLS NotificationFailureClass = "tls"
NotificationFailureConfiguration NotificationFailureClass = "configuration"
NotificationFailureRejected NotificationFailureClass = "rejected"
NotificationFailureUnknown NotificationFailureClass = "unknown"
)
// NotificationFailureClassCounts contains only bounded aggregate counters.
type NotificationFailureClassCounts struct {
Authentication int
RateLimited int
Connectivity int
TLS int
Configuration int
Rejected int
Unknown int
}
func (counts NotificationFailureClassCounts) AsMap() map[string]int {
return map[string]int{
string(NotificationFailureAuthentication): counts.Authentication,
string(NotificationFailureRateLimited): counts.RateLimited,
string(NotificationFailureConnectivity): counts.Connectivity,
string(NotificationFailureTLS): counts.TLS,
string(NotificationFailureConfiguration): counts.Configuration,
string(NotificationFailureRejected): counts.Rejected,
string(NotificationFailureUnknown): counts.Unknown,
}
}
// ClassifyNotificationFailure maps local error text to a fixed diagnostic
// bucket. It never returns any portion of the error string.
func ClassifyNotificationFailure(errorMessage string) NotificationFailureClass {
message := strings.ToLower(strings.TrimSpace(errorMessage))
containsAny := func(values ...string) bool {
for _, value := range values {
if strings.Contains(message, value) {
return true
}
}
return false
}
switch {
case containsAny(
"unauthorized", "forbidden", "authentication", "auth failed",
"auth negotiation",
"invalid token", "invalid api key", "invalid credentials",
"status 401", "http 401", "status 403", "http 403", "smtp 535", " 535 ",
):
return NotificationFailureAuthentication
case containsAny("rate limit", "too many requests", "status 429", "http 429"):
return NotificationFailureRateLimited
case containsAny(
"x509", "certificate", "tls", "ssl", "starttls",
):
return NotificationFailureTLS
case containsAny(
"timeout", "timed out", "deadline exceeded", "connection refused",
"connection reset", "network unreachable", "no such host", "dial tcp",
"broken pipe", "unexpected eof",
):
return NotificationFailureConnectivity
case containsAny(
"not configured", "missing ", "invalid address", "invalid url",
"invalid webhook", "invalid configuration",
"must start with http", "validation failed", "executable file not found",
"no supported", "requires payload template",
):
return NotificationFailureConfiguration
case containsAny(
"bad request", "not found", "method not allowed", "gone",
"payload too large", "unsupported media type", "unprocessable",
"status 400", "http 400", "status 404", "http 404",
"status 405", "http 405", "status 410", "http 410",
"status 413", "http 413", "status 415", "http 415",
"status 422", "http 422", "status 4", "http 4", "status 5", "http 5",
):
return NotificationFailureRejected
default:
return NotificationFailureUnknown
}
}
// QueuedNotification represents a notification in the persistent queue
type QueuedNotification struct {
ID string `json:"id"`
@@ -319,6 +411,7 @@ func (nq *NotificationQueue) initSchema() error {
attempts INTEGER,
success BOOLEAN,
error_message TEXT,
failure_class TEXT NOT NULL DEFAULT '',
payload_size INTEGER,
timestamp INTEGER NOT NULL,
FOREIGN KEY (notification_id) REFERENCES notification_queue(id)
@@ -353,9 +446,15 @@ func (nq *NotificationQueue) initSchema() error {
); err != nil {
return err
}
return nq.ensureJSONColumn(
if err := nq.ensureJSONColumn(
"notification_audit",
notificationOperationalLinksColumn,
); err != nil {
return err
}
return nq.ensureTextColumn(
"notification_audit",
notificationFailureClassColumn,
)
}
@@ -487,6 +586,22 @@ func (nq *NotificationQueue) ensureJSONColumn(table, column string) error {
return nil
}
func (nq *NotificationQueue) ensureTextColumn(table, column string) error {
columns, err := nq.tableColumns(table)
if err != nil {
return err
}
if columns[column] {
return nil
}
if _, err := nq.db.Exec(
`ALTER TABLE ` + table + ` ADD COLUMN ` + column + ` TEXT NOT NULL DEFAULT ''`,
); err != nil {
return fmt.Errorf("add %s.%s: %w", table, column, err)
}
return nil
}
func notificationLinksForAlerts(
alertsToLink []*alerts.Alert,
destinationID string,
@@ -1169,10 +1284,14 @@ func (nq *NotificationQueue) RecordAudit(notif *QueuedNotification, success bool
query := `
INSERT INTO notification_audit
(notification_id, type, method, status, alert_identifiers, alert_count, operational_links, attempts, success, error_message, payload_size, timestamp)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
(notification_id, type, method, status, alert_identifiers, alert_count, operational_links, attempts, success, error_message, failure_class, payload_size, timestamp)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`
failureClass := ""
if !success {
failureClass = string(ClassifyNotificationFailure(errorMsg))
}
_, err = nq.db.Exec(query,
notif.ID,
notif.Type,
@@ -1184,6 +1303,7 @@ func (nq *NotificationQueue) RecordAudit(notif *QueuedNotification, success bool
notif.Attempts,
success,
errorMsg,
failureClass,
notif.PayloadBytes,
time.Now().Unix(),
)
@@ -1249,7 +1369,8 @@ type TelemetryStats struct {
// Failures counts terminal delivery failures only. Recoverable failed
// attempts that returned to pending for retry are included in Attempts but
// not in Failures.
Failures int
Failures int
FailureClasses NotificationFailureClassCounts
}
// GetTelemetryStats returns aggregate delivery activity from locally retained
@@ -1276,10 +1397,33 @@ func (nq *NotificationQueue) GetTelemetryStats(since time.Time) (TelemetryStats,
COALESCE(SUM(CASE
WHEN success = 0 AND status IN ('failed', 'dlq') THEN 1
ELSE 0
END), 0),
COALESCE(SUM(CASE WHEN success = 0 AND status IN ('failed', 'dlq') AND failure_class = 'authentication' THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE WHEN success = 0 AND status IN ('failed', 'dlq') AND failure_class = 'rate_limited' THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE WHEN success = 0 AND status IN ('failed', 'dlq') AND failure_class = 'connectivity' THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE WHEN success = 0 AND status IN ('failed', 'dlq') AND failure_class = 'tls' THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE WHEN success = 0 AND status IN ('failed', 'dlq') AND failure_class = 'configuration' THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE WHEN success = 0 AND status IN ('failed', 'dlq') AND failure_class = 'rejected' THEN 1 ELSE 0 END), 0),
COALESCE(SUM(CASE
WHEN success = 0
AND status IN ('failed', 'dlq')
AND COALESCE(failure_class, '') NOT IN ('authentication', 'rate_limited', 'connectivity', 'tls', 'configuration', 'rejected')
THEN 1 ELSE 0
END), 0)
FROM notification_audit
WHERE timestamp >= ?
`, since.UTC().Unix()).Scan(&stats.Attempts, &stats.Deliveries, &stats.Failures)
`, since.UTC().Unix()).Scan(
&stats.Attempts,
&stats.Deliveries,
&stats.Failures,
&stats.FailureClasses.Authentication,
&stats.FailureClasses.RateLimited,
&stats.FailureClasses.Connectivity,
&stats.FailureClasses.TLS,
&stats.FailureClasses.Configuration,
&stats.FailureClasses.Rejected,
&stats.FailureClasses.Unknown,
)
if err != nil {
return TelemetryStats{}, fmt.Errorf("read notification telemetry aggregates: %w", err)
}
+79 -3
View File
@@ -524,9 +524,10 @@ func TestNewNotificationQueue_MigratesAuditAlertIdentifiersColumn(t *testing.T)
t.Fatalf("tableColumns(notification_queue) failed: %v", err)
}
if !queueColumns[notificationOperationalLinksColumn] ||
!columns[notificationOperationalLinksColumn] {
!columns[notificationOperationalLinksColumn] ||
!columns[notificationFailureClassColumn] {
t.Fatalf(
"expected migrated operational link columns, queue=%#v audit=%#v",
"expected migrated operational link and failure class columns, queue=%#v audit=%#v",
queueColumns,
columns,
)
@@ -1445,11 +1446,86 @@ func TestGetTelemetryStatsReturnsOnlyWindowedOutcomeCounts(t *testing.T) {
if err != nil {
t.Fatalf("GetTelemetryStats: %v", err)
}
if stats != (TelemetryStats{Attempts: 3, Deliveries: 1, Failures: 1}) {
if stats != (TelemetryStats{
Attempts: 3,
Deliveries: 1,
Failures: 1,
FailureClasses: NotificationFailureClassCounts{
Unknown: 1,
},
}) {
t.Fatalf("telemetry stats = %#v", stats)
}
}
func TestClassifyNotificationFailureUsesBoundedPrivacySafeBuckets(t *testing.T) {
tests := []struct {
errorMessage string
want NotificationFailureClass
}{
{"webhook returned status 401: Unauthorized", NotificationFailureAuthentication},
{"SMTP auth failed: 535 invalid credentials", NotificationFailureAuthentication},
{"webhook returned HTTP 429: Too Many Requests", NotificationFailureRateLimited},
{"x509: certificate signed by unknown authority", NotificationFailureTLS},
{"dial tcp 10.0.0.1:443: i/o timeout", NotificationFailureConnectivity},
{"apprise server URL is not configured", NotificationFailureConfiguration},
{"webhook returned status 413: Payload Too Large", NotificationFailureRejected},
{"provider-specific secret detail", NotificationFailureUnknown},
}
for _, test := range tests {
if got := ClassifyNotificationFailure(test.errorMessage); got != test.want {
t.Errorf("ClassifyNotificationFailure(%q) = %q, want %q", test.errorMessage, got, test.want)
}
}
}
func TestGetTelemetryStatsClassifiesTerminalFailuresOnly(t *testing.T) {
nq, err := NewNotificationQueue(t.TempDir())
if err != nil {
t.Fatalf("NewNotificationQueue: %v", err)
}
defer func() { _ = nq.Stop() }()
now := time.Now().UTC()
entries := []*QueuedNotification{
{ID: "auth-dlq", Type: "webhook", Status: QueueStatusDLQ, Attempts: 3, Config: []byte(`{}`), CreatedAt: now},
{ID: "rate-retry", Type: "webhook", Status: QueueStatusPending, Attempts: 1, Config: []byte(`{}`), CreatedAt: now},
{ID: "tls-failed", Type: "email", Status: QueueStatusFailed, Attempts: 3, Config: []byte(`{}`), CreatedAt: now},
}
for _, entry := range entries {
status := entry.Status
entry.Status = QueueStatusPending
if err := nq.Enqueue(entry); err != nil {
t.Fatalf("enqueue %s: %v", entry.ID, err)
}
entry.Status = status
}
if err := nq.RecordAudit(entries[0], false, "HTTP 401 Unauthorized"); err != nil {
t.Fatalf("record auth audit: %v", err)
}
if err := nq.RecordAudit(entries[1], false, "HTTP 429 Too Many Requests"); err != nil {
t.Fatalf("record retry audit: %v", err)
}
if err := nq.RecordAudit(entries[2], false, "x509 certificate failure"); err != nil {
t.Fatalf("record TLS audit: %v", err)
}
stats, err := nq.GetTelemetryStats(now.Add(-time.Hour))
if err != nil {
t.Fatalf("GetTelemetryStats: %v", err)
}
if stats.Attempts != 3 || stats.Failures != 2 {
t.Fatalf("telemetry totals = %#v, want 3 attempts and 2 terminal failures", stats)
}
if stats.FailureClasses.Authentication != 1 || stats.FailureClasses.TLS != 1 {
t.Fatalf("failure classes = %#v, want authentication=1 tls=1", stats.FailureClasses)
}
if stats.FailureClasses.RateLimited != 0 {
t.Fatalf("retry failure leaked into terminal class counts: %#v", stats.FailureClasses)
}
}
func TestPerformCleanup(t *testing.T) {
t.Run("cleanup removes old completed entries", func(t *testing.T) {
tempDir := t.TempDir()
+24 -2
View File
@@ -136,7 +136,8 @@ const (
// Schema v3 makes notification_failures_7d a terminal-delivery count;
// schema v4 adds complete approved-action outcome accounting, fixed
// pre-dispatch refusal categories, and verified finding-resolution linkage.
TelemetrySchemaVersion = 4
// Schema v5 adds bounded, content-free notification failure classes.
TelemetrySchemaVersion = 5
)
type installIDRecord struct {
@@ -237,7 +238,14 @@ type Ping struct {
NotificationDeliveries7d int `json:"notification_deliveries_7d"`
// NotificationFailures7d counts only terminal failed/dead-letter outcomes.
// Retry-attempt failures remain represented in NotificationAttempts7d.
NotificationFailures7d int `json:"notification_failures_7d"`
NotificationFailures7d int `json:"notification_failures_7d"`
NotificationFailuresAuthentication7d int `json:"notification_failures_authentication_7d"`
NotificationFailuresRateLimited7d int `json:"notification_failures_rate_limited_7d"`
NotificationFailuresConnectivity7d int `json:"notification_failures_connectivity_7d"`
NotificationFailuresTLS7d int `json:"notification_failures_tls_7d"`
NotificationFailuresConfiguration7d int `json:"notification_failures_configuration_7d"`
NotificationFailuresRejected7d int `json:"notification_failures_rejected_7d"`
NotificationFailuresUnknown7d int `json:"notification_failures_unknown_7d"`
// Pulse Intelligence usage (30-day counts/booleans — no prompts, commands, outputs, resource IDs, or token values)
PulseIntelligenceLoopConfigured bool `json:"pulse_intelligence_loop_configured"`
@@ -366,6 +374,13 @@ type Snapshot struct {
NotificationAttempts7d int
NotificationDeliveries7d int
NotificationFailures7d int
NotificationFailuresAuthentication7d int
NotificationFailuresRateLimited7d int
NotificationFailuresConnectivity7d int
NotificationFailuresTLS7d int
NotificationFailuresConfiguration7d int
NotificationFailuresRejected7d int
NotificationFailuresUnknown7d int
PulseIntelligenceLoopConfigured bool
PulseIntelligenceLoopActive30d bool
PulseIntelligenceCompleteOperationsLoop30d bool
@@ -930,6 +945,13 @@ func applySnapshot(base Ping, fn SnapshotFunc) Ping {
ping.NotificationAttempts7d = s.NotificationAttempts7d
ping.NotificationDeliveries7d = s.NotificationDeliveries7d // gitleaks:allow -- schema field name, not a credential
ping.NotificationFailures7d = s.NotificationFailures7d
ping.NotificationFailuresAuthentication7d = s.NotificationFailuresAuthentication7d
ping.NotificationFailuresRateLimited7d = s.NotificationFailuresRateLimited7d
ping.NotificationFailuresConnectivity7d = s.NotificationFailuresConnectivity7d
ping.NotificationFailuresTLS7d = s.NotificationFailuresTLS7d
ping.NotificationFailuresConfiguration7d = s.NotificationFailuresConfiguration7d
ping.NotificationFailuresRejected7d = s.NotificationFailuresRejected7d
ping.NotificationFailuresUnknown7d = s.NotificationFailuresUnknown7d
ping.PulseIntelligenceLoopConfigured = s.PulseIntelligenceLoopConfigured
ping.PulseIntelligenceLoopActive30d = s.PulseIntelligenceLoopActive30d
ping.PulseIntelligenceCompleteOperationsLoop30d = s.PulseIntelligenceCompleteOperationsLoop30d
+12 -4
View File
@@ -809,10 +809,13 @@ func TestBuildPreview_UsesCurrentHeartbeatPayload(t *testing.T) {
IsDocker: true,
GetSnapshot: func() Snapshot {
return Snapshot{
PVENodes: 3,
VMs: 10,
ActiveAlerts: 2,
AIEnabled: true,
PVENodes: 3,
VMs: 10,
ActiveAlerts: 2,
AIEnabled: true,
NotificationFailures7d: 3,
NotificationFailuresAuthentication7d: 2,
NotificationFailuresConnectivity7d: 1,
}
},
})
@@ -847,6 +850,11 @@ func TestBuildPreview_UsesCurrentHeartbeatPayload(t *testing.T) {
if preview.PVENodes != 3 || preview.VMs != 10 || preview.ActiveAlerts != 2 {
t.Fatalf("preview snapshot = %#v", preview)
}
if preview.NotificationFailures7d != 3 ||
preview.NotificationFailuresAuthentication7d != 2 ||
preview.NotificationFailuresConnectivity7d != 1 {
t.Fatalf("preview notification failure classes = %#v", preview)
}
if preview.InstallID == "" {
t.Fatal("expected preview install ID")
}
+7
View File
@@ -542,6 +542,13 @@ func Run(ctx context.Context, version string) error {
snap.NotificationAttempts7d = counts.NotificationAttempts7d
snap.NotificationDeliveries7d = counts.NotificationDeliveries7d
snap.NotificationFailures7d = counts.NotificationFailures7d
snap.NotificationFailuresAuthentication7d = counts.NotificationFailuresAuthentication7d
snap.NotificationFailuresRateLimited7d = counts.NotificationFailuresRateLimited7d
snap.NotificationFailuresConnectivity7d = counts.NotificationFailuresConnectivity7d
snap.NotificationFailuresTLS7d = counts.NotificationFailuresTLS7d
snap.NotificationFailuresConfiguration7d = counts.NotificationFailuresConfiguration7d
snap.NotificationFailuresRejected7d = counts.NotificationFailuresRejected7d
snap.NotificationFailuresUnknown7d = counts.NotificationFailuresUnknown7d
snap.DiscoveryEnabled = currentCfg.DiscoveryEnabled
// Feature flags from persisted config (using pre-created persistence).
+7
View File
@@ -92,6 +92,13 @@ USER_BASE_COUNT_FIELDS = (
("alerts_resolved_30d", "Alerts resolved (30d)"),
("notification_attempts_7d", "Notification attempts, including retries (7d)"),
("notification_deliveries_7d", "Notification deliveries (7d)"),
("notification_failures_authentication_7d", "Notification authentication failures (7d, schema v5+)"),
("notification_failures_rate_limited_7d", "Notification rate-limit failures (7d, schema v5+)"),
("notification_failures_connectivity_7d", "Notification connectivity failures (7d, schema v5+)"),
("notification_failures_tls_7d", "Notification TLS failures (7d, schema v5+)"),
("notification_failures_configuration_7d", "Notification configuration failures (7d, schema v5+)"),
("notification_failures_rejected_7d", "Notification destination rejections (7d, schema v5+)"),
("notification_failures_unknown_7d", "Notification unknown failures (7d, schema v5+)"),
)
NOTIFICATION_FAILURE_COUNT_SIGNALS = (
(