|
@@ -229,3 +289,24 @@ const statusOptions = [
+
+
diff --git a/frontend/src/components/__tests__/AlertList.spec.ts b/frontend/src/components/__tests__/AlertList.spec.ts
index 9c4371bf..bae1e814 100644
--- a/frontend/src/components/__tests__/AlertList.spec.ts
+++ b/frontend/src/components/__tests__/AlertList.spec.ts
@@ -8,6 +8,18 @@ import AlertList from '@/components/AlertList.vue'
import { useAlertsStore } from '@/stores/alerts'
import { detailSlideOverKey } from '@/composables/useDetailSlideOver'
import type { Alert } from '@/services/alertApi'
+import { LOCAL_AGENT } from '@/services/apiFetch'
+
+const sse = vi.hoisted(() => new Map void>())
+vi.mock('@/services/sseBus', () => ({
+ sseBus: {
+ on: (type: string, fn: (e: MessageEvent) => void) => sse.set(type, fn),
+ off: (type: string) => sse.delete(type),
+ connect: () => {},
+ disconnect: () => {},
+ connected: false,
+ },
+}))
function alertAt(id: string, firedAt: string): Alert {
return {
@@ -20,14 +32,132 @@ function alertAt(id: string, firedAt: string): Alert {
entity_type: 'container',
entity_id: id,
entity_name: id,
+ agent_id: LOCAL_AGENT,
fired_at: firedAt,
created_at: firedAt,
}
}
+function jsonResponse(body: unknown, status = 200): Response {
+ return new Response(JSON.stringify(body), { status, headers: { 'content-type': 'application/json' } })
+}
+
+function mountList(props: { linkedId?: string } = {}) {
+ return mount(AlertList, {
+ props,
+ attachTo: document.body,
+ global: {
+ provide: { [detailSlideOverKey as symbol]: { openDetail: vi.fn() } },
+ stubs: { AcknowledgeButton: true },
+ },
+ })
+}
+
describe('AlertList', () => {
afterEach(() => {
vi.unstubAllGlobals()
+ document.body.innerHTML = ''
+ })
+
+ it('highlights the linked alert when it is in the loaded page', async () => {
+ const fetchMock = vi.fn().mockImplementation(async () => jsonResponse(alertAt('a1', '2026-09-30T10:00:00Z')))
+ vi.stubGlobal('fetch', fetchMock)
+ setActivePinia(createPinia())
+ const store = useAlertsStore()
+ store.alerts = [alertAt('a2', '2026-09-30T11:00:00Z'), alertAt('a1', '2026-09-30T10:00:00Z')]
+
+ const wrapper = mountList({ linkedId: 'a1' })
+ await flushPromises()
+
+ const rows = wrapper.findAll('tr[data-alert-id]')
+ expect(rows.map((r) => r.attributes('data-alert-id'))).toEqual(['a2', 'a1'])
+ expect(rows[1]!.classes()).toContain('alert-linked')
+ expect(rows[0]!.classes()).not.toContain('alert-linked')
+ })
+
+ it('fetches the linked alert by id and shows it above a page that does not hold it', async () => {
+ const fetchMock = vi.fn().mockImplementation(async () => jsonResponse(alertAt('old', '2026-08-01T10:00:00Z')))
+ vi.stubGlobal('fetch', fetchMock)
+ setActivePinia(createPinia())
+ const store = useAlertsStore()
+ store.alerts = [alertAt('a2', '2026-09-30T11:00:00Z')]
+
+ const wrapper = mountList({ linkedId: 'old' })
+ await flushPromises()
+
+ expect(String(fetchMock.mock.calls[0]![0])).toMatch(/\/alerts\/old$/)
+ const rows = wrapper.findAll('tr[data-alert-id]')
+ expect(rows.map((r) => r.attributes('data-alert-id'))).toEqual(['old', 'a2'])
+ expect(rows[0]!.classes()).toContain('alert-linked')
+ })
+
+ it('applies live updates to a linked alert fetched outside the loaded page', async () => {
+ const old = alertAt('old', '2026-08-01T10:00:00Z')
+ vi.stubGlobal('fetch', vi.fn().mockImplementation(async () => jsonResponse(old)))
+ setActivePinia(createPinia())
+ const store = useAlertsStore()
+ store.alerts = [alertAt('a2', '2026-09-30T11:00:00Z')]
+ store.connectSSE()
+
+ const wrapper = mountList({ linkedId: 'old' })
+ await flushPromises()
+
+ sse.get('alert.acknowledged')!({ data: JSON.stringify({ ...old, acknowledged_at: '2026-10-01T09:00:00Z', acknowledged_by: 'ops' }) } as MessageEvent)
+ expect(store.alerts.find((a) => a.id === 'old')!.acknowledged_at).toBe('2026-10-01T09:00:00Z')
+
+ sse.get('alert.resolved')!({ data: JSON.stringify({ ...old, status: 'resolved', resolved_at: '2026-10-01T09:05:00Z' }) } as MessageEvent)
+ await flushPromises()
+ const row = wrapper.find('tr[data-alert-id="old"]')
+ expect(row.classes()).toContain('alert-linked')
+ expect(row.text()).toContain('resolved')
+ store.disconnectSSE()
+ })
+
+ it('updates a linked alert outside the loaded page when it is acknowledged from its row', async () => {
+ const old = alertAt('old', '2026-08-01T10:00:00Z')
+ const acked = { ...old, acknowledged_at: '2026-10-01T09:00:00Z', acknowledged_by: 'maintenant-ui' }
+ const fetchMock = vi.fn().mockImplementation(async (url: string) => jsonResponse(String(url).endsWith('/acknowledge') ? acked : old))
+ vi.stubGlobal('fetch', fetchMock)
+ setActivePinia(createPinia())
+ const store = useAlertsStore()
+ store.alerts = [alertAt('a2', '2026-09-30T11:00:00Z')]
+
+ mountList({ linkedId: 'old' })
+ await flushPromises()
+ await store.acknowledgeAlert('old')
+
+ expect(store.alerts.find((a) => a.id === 'old')!.acknowledged_at).toBe('2026-10-01T09:00:00Z')
+ })
+
+ it('keeps the linked alert when the first page arrives after it, and drops it on a filter change', async () => {
+ const old = alertAt('old', '2026-08-01T10:00:00Z')
+ const page = { alerts: [alertAt('a2', '2026-09-30T11:00:00Z')], has_more: false }
+ vi.stubGlobal('fetch', vi.fn().mockImplementation(async (url: string) => jsonResponse(String(url).includes('/alerts/old') ? old : page)))
+ setActivePinia(createPinia())
+ const store = useAlertsStore()
+
+ const wrapper = mountList({ linkedId: 'old' })
+ await flushPromises()
+ await store.fetchAlerts()
+ expect(store.alerts.map((a) => a.id)).toEqual(['old', 'a2'])
+
+ await wrapper.findAllComponents({ name: 'SelectInput' })[0]!.vm.$emit('update:modelValue', 'endpoint')
+ await flushPromises()
+ expect(store.alerts.map((a) => a.id)).toEqual(['a2'])
+ })
+
+ it('says so when the linked alert no longer exists', async () => {
+ const fetchMock = vi.fn().mockImplementation(async () =>
+ jsonResponse({ error: { code: 'NOT_FOUND', message: 'alert not found' } }, 404),
+ )
+ vi.stubGlobal('fetch', fetchMock)
+ setActivePinia(createPinia())
+
+ const wrapper = mountList({ linkedId: 'gone' })
+ await flushPromises()
+
+ expect(wrapper.text()).toContain('Linked alert not found')
+ expect(wrapper.text()).toContain('This alert no longer exists')
})
it('pages from the last alert listed, not only from its second', async () => {
diff --git a/frontend/src/pages/AlertsPage.vue b/frontend/src/pages/AlertsPage.vue
index f69760a1..952f7cdf 100644
--- a/frontend/src/pages/AlertsPage.vue
+++ b/frontend/src/pages/AlertsPage.vue
@@ -23,8 +23,14 @@ const router = useRouter()
const store = useAlertsStore()
const triggersStore = useTriggersStore()
+const linkedAlertId = computed(() => {
+ const id = route.query.alert
+ return typeof id === 'string' && id !== '' ? id : undefined
+})
+
const activeTab = computed({
get: () => {
+ if (linkedAlertId.value) return 'history'
const t = route.params.tab as string
if (t === 'channels') return 'triggers' // legacy redirect
if (t === 'triggers' || t === 'silence' || t === 'history') return t
@@ -77,7 +83,7 @@ onUnmounted(() => {
-
+
resolved_by_id?: string | null
acknowledged_at?: string | null
diff --git a/frontend/src/stores/alerts.ts b/frontend/src/stores/alerts.ts
index 4ed877e9..c9583bf9 100644
--- a/frontend/src/stores/alerts.ts
+++ b/frontend/src/stores/alerts.ts
@@ -27,6 +27,7 @@ export const useAlertsStore = defineStore('alerts', () => {
info: [],
})
const hasMore = ref(false)
+ const pinnedId = ref(null)
const loading = ref(false)
const error = ref(null)
@@ -51,7 +52,8 @@ export const useAlertsStore = defineStore('alerts', () => {
if (params?.before) {
alerts.value = [...alerts.value, ...res.alerts]
} else {
- alerts.value = res.alerts
+ const pinned = alerts.value.find((a) => a.id === pinnedId.value)
+ alerts.value = pinned && !res.alerts.some((a) => a.id === pinned.id) ? [pinned, ...res.alerts] : res.alerts
}
hasMore.value = res.has_more
} catch (e) {
@@ -227,6 +229,17 @@ export const useAlertsStore = defineStore('alerts', () => {
}
}
+ function pinAlert(alert: Alert) {
+ pinnedId.value = alert.id
+ if (!alerts.value.some((a) => a.id === alert.id)) {
+ alerts.value = [alert, ...alerts.value]
+ }
+ }
+
+ function unpinAlert() {
+ pinnedId.value = null
+ }
+
function updateAlertInList(updated: Alert) {
const idx = alerts.value.findIndex((a) => a.id === updated.id)
if (idx >= 0) {
@@ -249,6 +262,8 @@ export const useAlertsStore = defineStore('alerts', () => {
fetchActiveAlerts,
fetchSilenceRules,
acknowledgeAlert,
+ pinAlert,
+ unpinAlert,
clearNewAlertCount,
connectSSE,
disconnectSSE,
diff --git a/internal/alert/down.go b/internal/alert/down.go
index 321d6346..f47e4a75 100644
--- a/internal/alert/down.go
+++ b/internal/alert/down.go
@@ -94,6 +94,7 @@ func (d *DownDetector) Check(ctx context.Context) error {
EntityType: "container",
EntityID: c.ID,
EntityName: c.Name,
+ AgentID: c.AgentID,
Details: map[string]any{
"state": string(c.State),
"stopped_for": int64(stoppedFor / time.Second),
@@ -122,6 +123,7 @@ func (d *DownDetector) Check(ctx context.Context) error {
EntityType: "container",
EntityID: id,
EntityName: name,
+ AgentID: a.AgentID,
Timestamp: now,
})
}
diff --git a/internal/alert/down_test.go b/internal/alert/down_test.go
index e0521719..497b27dd 100644
--- a/internal/alert/down_test.go
+++ b/internal/alert/down_test.go
@@ -17,6 +17,7 @@ import (
"github.com/kolapsis/maintenant/internal/container"
"github.com/kolapsis/maintenant/internal/store"
"github.com/kolapsis/maintenant/internal/store/storetest"
+ "github.com/kolapsis/maintenant/internal/uid"
)
const downThreshold = 5 * time.Minute
@@ -73,6 +74,7 @@ func (f *downFixture) fire(t *testing.T, evt alert.Event) {
EntityType: evt.EntityType,
EntityID: evt.EntityID,
EntityName: evt.EntityName,
+ AgentID: evt.AgentID,
Details: "{}",
FiredAt: evt.Timestamp,
})
@@ -100,6 +102,7 @@ func TestDownDetector_FiresPastTheThreshold(t *testing.T) {
assert.Equal(t, "container", evt.EntityType)
assert.Equal(t, id, evt.EntityID)
assert.Equal(t, "api", evt.EntityName)
+ assert.Equal(t, uid.LocalAgent, evt.AgentID)
assert.Equal(t, int64(downThreshold/time.Second), evt.Details["threshold_seconds"])
}
@@ -164,6 +167,7 @@ func TestDownDetector_ResolvesWhenTheContainerRunsAgain(t *testing.T) {
assert.Equal(t, alert.SeverityInfo, events[0].Severity)
assert.Equal(t, alert.AlertTypeContainerDown, events[0].AlertType)
assert.Equal(t, "api", events[0].EntityName)
+ assert.Equal(t, uid.LocalAgent, events[0].AgentID)
}
// Nothing is emitted for a container that never crossed the threshold, so a
diff --git a/internal/alert/engine.go b/internal/alert/engine.go
index 5d7d4442..1ca29c1b 100644
--- a/internal/alert/engine.go
+++ b/internal/alert/engine.go
@@ -13,6 +13,7 @@ import (
"time"
"github.com/kolapsis/maintenant/internal/event"
+ "github.com/kolapsis/maintenant/internal/uid"
)
const engineChannelBuffer = 256
@@ -272,6 +273,7 @@ func (e *Engine) processEvent(ctx context.Context, evt Event) {
EntityType: evt.EntityType,
EntityID: evt.EntityID,
EntityName: evt.EntityName,
+ AgentID: uid.Agent(evt.AgentID),
Details: detailsJSON,
FiredAt: evt.Timestamp,
}
@@ -480,6 +482,7 @@ func (e *Engine) processRecovery(ctx context.Context, evt Event) {
EntityType: evt.EntityType,
EntityID: evt.EntityID,
EntityName: evt.EntityName,
+ AgentID: activeAlert.AgentID,
Details: detailsJSON,
FiredAt: evt.Timestamp,
}
@@ -739,6 +742,7 @@ func alertToMap(a *Alert) map[string]interface{} {
"entity_type": a.EntityType,
"entity_id": a.EntityID,
"entity_name": a.EntityName,
+ "agent_id": a.AgentID,
"fired_at": a.FiredAt.UTC().Format(time.RFC3339),
"created_at": a.CreatedAt.UTC().Format(time.RFC3339),
}
diff --git a/internal/alert/engine_agent_test.go b/internal/alert/engine_agent_test.go
new file mode 100644
index 00000000..f71b7c47
--- /dev/null
+++ b/internal/alert/engine_agent_test.go
@@ -0,0 +1,96 @@
+// Copyright 2026 Benjamin Touchard (Kolapsis)
+// SPDX-License-Identifier: Apache-2.0
+
+package alert_test
+
+import (
+ "context"
+ "io"
+ "log/slog"
+ "testing"
+ "time"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+
+ "github.com/kolapsis/maintenant/internal/alert"
+ "github.com/kolapsis/maintenant/internal/store"
+ "github.com/kolapsis/maintenant/internal/store/storetest"
+ "github.com/kolapsis/maintenant/internal/uid"
+)
+
+func TestEngine_AgentIDFollowsTheAlert(t *testing.T) {
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
+
+ logger := slog.New(slog.NewTextHandler(io.Discard, nil))
+ db := storetest.Open(t, logger)
+ alertStore := store.NewAlertStore(db)
+ eng := alert.NewEngine(alert.EngineDeps{
+ AlertStore: alertStore,
+ ChannelStore: store.NewChannelStore(db),
+ TriggerStore: store.NewTriggerStore(db),
+ SilenceStore: store.NewSilenceStore(db),
+ Logger: logger,
+ })
+ eng.Start(ctx)
+
+ const remote = "0190a1b2-c3d4-7e5f-8a9b-0c1d2e3f4a5b"
+ event := func(entityID, agentID, severity string, recover bool) alert.Event {
+ return alert.Event{
+ Source: alert.SourceContainer,
+ AlertType: "health_unhealthy",
+ Severity: severity,
+ IsRecover: recover,
+ Message: "unhealthy",
+ EntityType: "container",
+ EntityID: entityID,
+ EntityName: entityID,
+ AgentID: agentID,
+ Timestamp: time.Now(),
+ }
+ }
+ activeOf := func(entityID string) *alert.Alert {
+ active, err := alertStore.ListActiveAlerts(ctx)
+ if err != nil {
+ return nil
+ }
+ for _, a := range active {
+ if a.EntityID == entityID {
+ return a
+ }
+ }
+ return nil
+ }
+
+ eng.EventChannel() <- event("remote-web", remote, alert.SeverityWarning, false)
+ eng.EventChannel() <- event("local-web", "", alert.SeverityWarning, false)
+ require.Eventually(t, func() bool { return activeOf("remote-web") != nil && activeOf("local-web") != nil },
+ 5*time.Second, 10*time.Millisecond)
+ assert.Equal(t, remote, activeOf("remote-web").AgentID)
+ assert.Equal(t, uid.LocalAgent, activeOf("local-web").AgentID, "an event without an agent belongs to the local runtime")
+
+ eng.EventChannel() <- event("remote-web", "", alert.SeverityCritical, false)
+ require.Eventually(t, func() bool {
+ a := activeOf("remote-web")
+ return a != nil && a.Severity == alert.SeverityCritical
+ }, 5*time.Second, 10*time.Millisecond)
+ assert.Equal(t, remote, activeOf("remote-web").AgentID, "a severity raise keeps the agent")
+
+ original := activeOf("remote-web")
+ eng.EventChannel() <- event("remote-web", "", alert.SeverityInfo, true)
+ var resolved *alert.Alert
+ require.Eventually(t, func() bool {
+ a, err := alertStore.GetAlert(ctx, original.ID)
+ if err != nil || a == nil || a.ResolvedByID == nil {
+ return false
+ }
+ resolved = a
+ return true
+ }, 5*time.Second, 10*time.Millisecond)
+
+ recovery, err := alertStore.GetAlert(ctx, *resolved.ResolvedByID)
+ require.NoError(t, err)
+ require.NotNil(t, recovery)
+ assert.Equal(t, remote, recovery.AgentID, "the recovery record inherits the agent of the alert it resolves")
+}
diff --git a/internal/alert/model.go b/internal/alert/model.go
index c7165835..ecb71f41 100644
--- a/internal/alert/model.go
+++ b/internal/alert/model.go
@@ -58,6 +58,7 @@ type Event struct {
EntityType string // "container", "endpoint", "heartbeat", "certificate"
EntityID string // UUID of the referenced entity in its source table
EntityName string // display name
+ AgentID string // host the entity belongs to; empty means the local runtime
Details map[string]any // source-specific metadata
Timestamp time.Time // when condition was detected
}
@@ -73,6 +74,7 @@ type Alert struct {
EntityType string `json:"entity_type"`
EntityID string `json:"entity_id"`
EntityName string `json:"entity_name"`
+ AgentID string `json:"agent_id"`
Details string `json:"details"`
ResolvedByID *string `json:"resolved_by_id"`
FiredAt time.Time `json:"fired_at"`
diff --git a/internal/alert/notifier.go b/internal/alert/notifier.go
index 1c174735..76877206 100644
--- a/internal/alert/notifier.go
+++ b/internal/alert/notifier.go
@@ -9,13 +9,16 @@ import (
"fmt"
"hash/fnv"
"log/slog"
+ "net"
"net/http"
+ "net/url"
"strings"
"sync"
"time"
"github.com/kolapsis/maintenant/internal/event"
"github.com/kolapsis/maintenant/internal/ssrf"
+ "github.com/kolapsis/maintenant/internal/uid"
)
const (
@@ -68,28 +71,41 @@ type Notifier struct {
logger *slog.Logger
webhook *webhookSender
senders map[string]ChannelSender
+ baseURL string
suspendedMu sync.Mutex
suspendedLogged map[string]bool
}
+// NotifierOption configures a Notifier at construction.
+type NotifierOption func(*Notifier)
+
+// WithBaseURL sets the public URL of the UI that webhook payloads link each alert to.
+func WithBaseURL(baseURL string) NotifierOption {
+ return func(n *Notifier) {
+ n.baseURL = usableBaseURL(baseURL)
+ }
+}
+
// NewNotifier creates a new webhook notifier. Its HTTP client blocks delivery
// to private/internal IPs (SSRF guard) at dial time unless allowPrivate is set
// (dev only, via MAINTENANT_ALLOW_PRIVATE_WEBHOOKS).
-func NewNotifier(channelStore ChannelStore, logger *slog.Logger, allowPrivate bool) *Notifier {
+func NewNotifier(channelStore ChannelStore, logger *slog.Logger, allowPrivate bool, opts ...NotifierOption) *Notifier {
client := ssrf.NewHTTPClient(webhookTimeout, allowPrivate)
- webhook := &webhookSender{client: client, format: formatWebhookPayload, logger: logger}
n := &Notifier{
- channelStore: channelStore,
- httpClient: client,
- logger: logger,
- webhook: webhook,
- senders: map[string]ChannelSender{
- "webhook": webhook,
- "discord": NewWebhookSender(client, formatDiscordPayload, logger),
- },
+ channelStore: channelStore,
+ httpClient: client,
+ logger: logger,
suspendedLogged: make(map[string]bool),
}
+ for _, opt := range opts {
+ opt(n)
+ }
+ n.webhook = &webhookSender{client: client, format: n.formatWebhookPayload, logger: logger}
+ n.senders = map[string]ChannelSender{
+ "webhook": n.webhook,
+ "discord": NewWebhookSender(client, formatDiscordPayload, logger),
+ }
for i := range n.queues {
n.queues[i] = make(chan NotificationJob, notifierChannelBuffer)
}
@@ -324,6 +340,7 @@ func (n *Notifier) SendTestWebhook(ctx context.Context, ch *NotificationChannel)
Message: "maintenant test notification",
EntityType: "test",
EntityName: "test",
+ AgentID: uid.LocalAgent,
FiredAt: time.Now().UTC(),
CreatedAt: time.Now().UTC(),
}
@@ -338,15 +355,35 @@ func (n *Notifier) SendTestWebhook(ctx context.Context, ch *NotificationChannel)
return sender.SendTest(ctx, ch, testAlert)
}
-func formatWebhookPayload(eventType string, a *Alert) ([]byte, error) {
+func (n *Notifier) formatWebhookPayload(eventType string, a *Alert) ([]byte, error) {
+ m := alertToMap(a)
+ if n.baseURL != "" && a.ID != "" {
+ m["url"] = n.baseURL + "/alerts/history?alert=" + url.QueryEscape(a.ID)
+ }
payload := WebhookPayload{
Event: eventType,
- Alert: alertToMap(a),
+ Alert: m,
Timestamp: time.Now().UTC().Format(time.RFC3339),
}
return json.Marshal(payload)
}
+// usableBaseURL returns baseURL without its trailing slash, or "" when it names no host a recipient could open.
+func usableBaseURL(baseURL string) string {
+ u, err := url.Parse(strings.TrimRight(baseURL, "/"))
+ if err != nil || (u.Scheme != "http" && u.Scheme != "https") {
+ return ""
+ }
+ host := u.Hostname()
+ if host == "" {
+ return ""
+ }
+ if ip := net.ParseIP(host); ip != nil && ip.IsUnspecified() {
+ return ""
+ }
+ return u.String()
+}
+
func formatDiscordPayload(eventType string, a *Alert) ([]byte, error) {
color := severityColor(a.Severity)
diff --git a/internal/alert/notifier_test.go b/internal/alert/notifier_test.go
index 86684486..a4dd9b55 100644
--- a/internal/alert/notifier_test.go
+++ b/internal/alert/notifier_test.go
@@ -327,3 +327,42 @@ func TestSendTestWebhook_UnregisteredTypeIsAGenericWebhook(t *testing.T) {
_, ok := n.Validator("webhook")
assert.False(t, ok)
}
+
+func TestWebhookPayload_CarriesAgentAndLink(t *testing.T) {
+ const alertID = "0190a1b2-0000-7000-8000-000000000001"
+ const remote = "0190a1b2-c3d4-7e5f-8a9b-0c1d2e3f4a5b"
+
+ tests := []struct {
+ name string
+ opts []NotifierOption
+ wantURL string
+ }{
+ {"public base url", []NotifierOption{WithBaseURL("https://mnt.example.com/")}, "https://mnt.example.com/alerts/history?alert=" + alertID},
+ {"base url under a path", []NotifierOption{WithBaseURL("https://example.com/maintenant")}, "https://example.com/maintenant/alerts/history?alert=" + alertID},
+ {"listen address without host", []NotifierOption{WithBaseURL("http://:8080")}, ""},
+ {"unspecified address", []NotifierOption{WithBaseURL("http://0.0.0.0:8080")}, ""},
+ {"no base url", nil, ""},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ srv, body, _ := captureServer(t, http.StatusOK)
+ n := NewNotifier(nil, slog.New(slog.NewTextHandler(io.Discard, nil)), true, tt.opts...)
+ a := &Alert{
+ ID: alertID, Source: SourceContainer, AlertType: "health_unhealthy", Severity: SeverityWarning,
+ Status: StatusActive, Message: "unhealthy", EntityType: "container", EntityID: "c1", EntityName: "web",
+ AgentID: remote, FiredAt: time.Now(), CreatedAt: time.Now(),
+ }
+ require.NoError(t, n.SendNow(context.Background(), a, &NotificationChannel{Type: "webhook", URL: srv.URL}))
+
+ var payload WebhookPayload
+ require.NoError(t, json.Unmarshal(*body, &payload))
+ assert.Equal(t, remote, payload.Alert["agent_id"])
+ link, ok := payload.Alert["url"]
+ if tt.wantURL == "" {
+ assert.False(t, ok, "url must be absent, got %v", link)
+ return
+ }
+ assert.Equal(t, tt.wantURL, link)
+ })
+ }
+}
diff --git a/internal/app/app.go b/internal/app/app.go
index 3c8881c8..eff04141 100644
--- a/internal/app/app.go
+++ b/internal/app/app.go
@@ -369,7 +369,7 @@ func New(cfg Config, logger *slog.Logger, opts ...Option) (*App, error) {
Password: cfg.SMTP.Password,
From: cfg.SMTP.From,
}
- a.notifier = alert.NewNotifier(channelStore, logger, cfg.AllowPrivateWebhooks)
+ a.notifier = alert.NewNotifier(channelStore, logger, cfg.AllowPrivateWebhooks, alert.WithBaseURL(cfg.BaseURL))
if a.ext.Channels != nil {
channels := a.ext.Channels(extpoint.ChannelDeps{
HTTPClient: a.notifier.HTTPClient(),
diff --git a/internal/app/kubernetes_alerts_test.go b/internal/app/kubernetes_alerts_test.go
index 5d462f6e..c9eca634 100644
--- a/internal/app/kubernetes_alerts_test.go
+++ b/internal/app/kubernetes_alerts_test.go
@@ -16,6 +16,7 @@ import (
v1 "github.com/kolapsis/maintenant/internal/api/v1"
"github.com/kolapsis/maintenant/internal/event"
"github.com/kolapsis/maintenant/internal/kubernetes"
+ "github.com/kolapsis/maintenant/internal/uid"
)
type fakeCluster struct {
@@ -100,6 +101,7 @@ func TestKubernetesReconcile_RaisesAndResolvesAlerts(t *testing.T) {
assert.Equal(t, "pod", fired.EntityType)
assert.Equal(t, "shop/api-x", fired.EntityID)
assert.Equal(t, alert.SeverityCritical, fired.Severity)
+ assert.Equal(t, uid.LocalAgent, fired.AgentID)
seen := map[string]bool{}
require.Eventually(t, func() bool {
diff --git a/internal/app/wiring.go b/internal/app/wiring.go
index 34382d05..5235da99 100644
--- a/internal/app/wiring.go
+++ b/internal/app/wiring.go
@@ -17,6 +17,7 @@ import (
"github.com/kolapsis/maintenant/internal/extension"
"github.com/kolapsis/maintenant/internal/heartbeat"
"github.com/kolapsis/maintenant/internal/security"
+ "github.com/kolapsis/maintenant/internal/uid"
"github.com/kolapsis/maintenant/internal/update"
)
@@ -35,6 +36,7 @@ func heartbeatAlertEvents(h *heartbeat.Heartbeat, alertType string, details map[
EntityType: "heartbeat",
EntityID: h.ID,
EntityName: h.Name,
+ AgentID: h.AgentID,
Details: details,
}
}
@@ -80,6 +82,7 @@ func certificateRecoveryEvent(m map[string]any) alert.Event {
EntityType: "certificate",
EntityID: toString(m["monitor_id"]),
EntityName: host,
+ AgentID: agentIDOf(m),
Details: m,
}
}
@@ -124,6 +127,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: ra.ContainerID,
EntityName: ra.ContainerName,
+ AgentID: ra.AgentID,
Details: map[string]any{
"restart_count": ra.RestartCount,
"threshold": ra.Threshold,
@@ -142,6 +146,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: toString(m["container_id"]),
EntityName: toString(m["container_name"]),
+ AgentID: agentIDOf(m),
Timestamp: time.Now(),
})
}
@@ -169,6 +174,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: toString(m["id"]),
EntityName: toString(m["container_name"]),
+ AgentID: agentIDOf(m),
Details: m,
Timestamp: time.Now(),
})
@@ -182,6 +188,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: toString(m["id"]),
EntityName: toString(m["container_name"]),
+ AgentID: agentIDOf(m),
Details: m,
Timestamp: time.Now(),
})
@@ -218,6 +225,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "endpoint",
EntityID: ep.ID,
EntityName: epName,
+ AgentID: ep.AgentID,
Details: map[string]any{
"target": ep.Target,
"reason": result.DegradedReason,
@@ -234,6 +242,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "endpoint",
EntityID: ep.ID,
EntityName: epName,
+ AgentID: ep.AgentID,
Timestamp: result.Timestamp,
})
}
@@ -255,6 +264,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "endpoint",
EntityID: al.EndpointID,
EntityName: entityName,
+ AgentID: ep.AgentID,
Details: map[string]any{
"target": al.Target,
"failures": al.Failures,
@@ -282,6 +292,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "endpoint",
EntityID: al.EndpointID,
EntityName: entityName,
+ AgentID: ep.AgentID,
Details: map[string]any{
"target": al.Target,
"successes": al.Successes,
@@ -374,6 +385,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "certificate",
EntityID: toString(m["monitor_id"]),
EntityName: toString(m["hostname"]),
+ AgentID: agentIDOf(m),
Details: m,
Timestamp: time.Now(),
})
@@ -402,6 +414,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: toString(m["container_id"]),
EntityName: toString(m["container_name"]),
+ AgentID: agentIDOf(m),
Details: m,
Timestamp: time.Now(),
})
@@ -416,6 +429,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: toString(m["container_id"]),
EntityName: toString(m["container_name"]),
+ AgentID: agentIDOf(m),
Details: m,
Timestamp: time.Now(),
})
@@ -424,6 +438,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
// Security insight alerts
a.securitySvc.SetAlertCallback(func(containerID string, containerName string, insights []security.Insight, isRecover bool) {
+ agentID := a.containerAgentID(ctx, containerID)
if isRecover {
sendAlert(alert.Event{
Source: alert.SourceSecurity,
@@ -434,6 +449,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: containerID,
EntityName: containerName,
+ AgentID: agentID,
Details: map[string]any{},
Timestamp: time.Now(),
})
@@ -448,6 +464,7 @@ func (a *App) wireAlertCallbacks(alertDetector *alert.EndpointAlertDetector) {
EntityType: "container",
EntityID: containerID,
EntityName: containerName,
+ AgentID: agentID,
Details: map[string]any{
"insight_count": fmt.Sprintf("%d", len(insights)),
"highest_severity": hs,
@@ -507,6 +524,7 @@ func updateDetectedAlert(m map[string]any, withChangelog bool) alert.Event {
EntityType: "container",
EntityID: entityID,
EntityName: containerName,
+ AgentID: agentIDOf(m),
Details: details,
}
}
@@ -543,6 +561,7 @@ func updateResolvedAlert(m map[string]any) alert.Event {
EntityType: "container",
EntityID: entityID,
EntityName: containerName,
+ AgentID: agentIDOf(m),
}
}
@@ -609,6 +628,7 @@ func (a *App) wirePostureCallbacks() {
EntityType: "infrastructure",
EntityID: "",
EntityName: "infrastructure",
+ AgentID: uid.LocalAgent,
Details: map[string]any{
"score": score,
"previous_score": previousScore,
@@ -634,19 +654,19 @@ func (a *App) wireSwarmCallbacks(m *swarmManager) {
m.events.SetNodeService(m.nodeSvc)
m.nodeSvc.SetEventCallback(sseBroadcast)
- m.nodeSvc.SetAlertCallback(a.emitAlert)
+ m.nodeSvc.SetAlertCallback(a.emitLocalRuntimeAlert)
m.crashLoop.SetEventCallback(sseBroadcast)
- m.crashLoop.SetAlertCallback(a.emitAlert)
+ m.crashLoop.SetAlertCallback(a.emitLocalRuntimeAlert)
m.updateTracker.SetEventCallback(sseBroadcast)
- m.updateTracker.SetAlertCallback(a.emitAlert)
+ m.updateTracker.SetAlertCallback(a.emitLocalRuntimeAlert)
m.replicaChecker.SetEventCallback(sseBroadcast)
- m.replicaChecker.SetAlertCallback(a.emitAlert)
+ m.replicaChecker.SetAlertCallback(a.emitLocalRuntimeAlert)
}
// wireKubernetesAlerts routes the local cluster's alerts to the alert engine and
// their SSE events to the broker.
func (a *App) wireKubernetesAlerts() {
- a.k8sAlerts.SetAlertCallback(a.emitAlert)
+ a.k8sAlerts.SetAlertCallback(a.emitLocalRuntimeAlert)
a.k8sAlerts.SetEventCallback(func(eventType string, data any) {
a.broker.Broadcast(v1.SSEEvent{Type: eventType, Data: data})
})
@@ -665,6 +685,7 @@ func agentLifecycleEvent(agentID, name, reason string, connected bool) alert.Eve
EntityType: "agent",
EntityID: agentID,
EntityName: name,
+ AgentID: agentID,
}
switch {
case connected:
@@ -726,6 +747,21 @@ func (a *App) emitAlert(evt alert.Event) {
a.statusSvc.HandleAlertEvent(context.Background(), evt)
}
+// emitLocalRuntimeAlert pushes an alert raised by the server's own Swarm or Kubernetes runtime.
+func (a *App) emitLocalRuntimeAlert(evt alert.Event) {
+ evt.AgentID = uid.LocalAgent
+ a.emitAlert(evt)
+}
+
+// containerAgentID returns the agent that runs the container, empty when the container is unknown.
+func (a *App) containerAgentID(ctx context.Context, containerID string) string {
+ c, err := a.containerSvc.GetContainer(ctx, containerID)
+ if err != nil || c == nil {
+ return ""
+ }
+ return c.AgentID
+}
+
// broadcastAgentUpdated republishes an agent over SSE after its identity changed.
func (a *App) broadcastAgentUpdated(ctx context.Context, agentID string) {
ag, err := a.agentStore.Get(ctx, agentID)
@@ -764,6 +800,11 @@ func (a *App) startOSEOL(ctx context.Context) {
}()
}
+func agentIDOf(m map[string]any) string {
+ id, _ := m["agent_id"].(string)
+ return id
+}
+
func toString(v any) string {
if s, ok := v.(string); ok {
return s
diff --git a/internal/app/wiring_agent_lifecycle_test.go b/internal/app/wiring_agent_lifecycle_test.go
index 705d8cd7..10108038 100644
--- a/internal/app/wiring_agent_lifecycle_test.go
+++ b/internal/app/wiring_agent_lifecycle_test.go
@@ -42,6 +42,9 @@ func TestAgentLifecycleEvent(t *testing.T) {
if evt.EntityID != "agent-uuid" {
t.Errorf("entityID = %q, want agent-uuid", evt.EntityID)
}
+ if evt.AgentID != "agent-uuid" {
+ t.Errorf("agentID = %q, want agent-uuid", evt.AgentID)
+ }
if evt.Message == "" {
t.Error("message must not be empty")
}
diff --git a/internal/app/wiring_alert_entity_test.go b/internal/app/wiring_alert_entity_test.go
index e6f9163d..0c6458e9 100644
--- a/internal/app/wiring_alert_entity_test.go
+++ b/internal/app/wiring_alert_entity_test.go
@@ -22,6 +22,15 @@ import (
// An event without an entity name reaches every channel titled "Alert: ".
func TestAlertEvents_BuiltInThisPackageNameTheirEntity(t *testing.T) {
+ assertEveryAlertEventSets(t, "EntityName")
+}
+
+func TestAlertEvents_BuiltInThisPackageNameTheirAgent(t *testing.T) {
+ assertEveryAlertEventSets(t, "AgentID")
+}
+
+func assertEveryAlertEventSets(t *testing.T, field string) {
+ t.Helper()
files, err := filepath.Glob("*.go")
require.NoError(t, err)
@@ -39,7 +48,7 @@ func TestAlertEvents_BuiltInThisPackageNameTheirEntity(t *testing.T) {
return true
}
built++
- assert.True(t, setsField(lit, "EntityName"), "%s builds an alert.Event without EntityName", fset.Position(lit.Pos()))
+ assert.True(t, setsField(lit, field), "%s builds an alert.Event without %s", fset.Position(lit.Pos()), field)
return true
})
}
diff --git a/internal/app/wiring_certificate_test.go b/internal/app/wiring_certificate_test.go
index 09bd5ecd..24d062bc 100644
--- a/internal/app/wiring_certificate_test.go
+++ b/internal/app/wiring_certificate_test.go
@@ -24,12 +24,14 @@ func TestCertificateRecoveryEvent_ResolvesTheRecoveredAlertType(t *testing.T) {
"monitor_id": "cert-1",
"hostname": "app.example.com",
"previous_alert_type": alertType,
+ "agent_id": "agent-1",
})
assert.Equal(t, alertType, evt.AlertType, "the recovery must carry the dedup key of the alert it resolves")
assert.True(t, evt.IsRecover)
assert.Equal(t, alert.SourceCertificate, evt.Source)
assert.Equal(t, "certificate", evt.EntityType)
assert.Equal(t, "cert-1", evt.EntityID)
+ assert.Equal(t, "agent-1", evt.AgentID)
assert.Contains(t, evt.Message, "app.example.com")
}
}
diff --git a/internal/app/wiring_heartbeat_test.go b/internal/app/wiring_heartbeat_test.go
index c589c1c9..3bb9b931 100644
--- a/internal/app/wiring_heartbeat_test.go
+++ b/internal/app/wiring_heartbeat_test.go
@@ -17,7 +17,7 @@ import (
// recovery must clear both dedup keys a failing ping can have used (deadline_missed
// and exit_code_failure), since the engine resolves by exact key.
func TestHeartbeatAlertEvents_Recovery_ClearsBothAlertTypes(t *testing.T) {
- h := &heartbeat.Heartbeat{ID: "hb-1", Name: "backup"}
+ h := &heartbeat.Heartbeat{ID: "hb-1", Name: "backup", AgentID: "agent-1"}
events := heartbeatAlertEvents(h, "recovery", map[string]any{"heartbeat_id": "hb-1"})
@@ -31,6 +31,7 @@ func TestHeartbeatAlertEvents_Recovery_ClearsBothAlertTypes(t *testing.T) {
assert.Equal(t, alert.SourceHeartbeat, evt.Source)
assert.Equal(t, "heartbeat", evt.EntityType)
assert.Equal(t, "hb-1", evt.EntityID)
+ assert.Equal(t, "agent-1", evt.AgentID)
}
assert.True(t, types[heartbeat.AlertTypeDeadlineMissed])
assert.True(t, types[heartbeat.AlertTypeExitCodeFailure])
diff --git a/internal/app/wiring_update_test.go b/internal/app/wiring_update_test.go
index d870bce1..f32c7971 100644
--- a/internal/app/wiring_update_test.go
+++ b/internal/app/wiring_update_test.go
@@ -142,3 +142,14 @@ func TestUpdateResolvedAlert(t *testing.T) {
t.Errorf("message = %q", evt.Message)
}
}
+
+func TestUpdateAlerts_CarryTheContainerAgent(t *testing.T) {
+ p := detectedPayload("svc", "id", 90)
+ p["agent_id"] = "agent-1"
+ if got := updateDetectedAlert(p, false).AgentID; got != "agent-1" {
+ t.Errorf("detected AgentID = %q, want agent-1", got)
+ }
+ if got := updateResolvedAlert(p).AgentID; got != "agent-1" {
+ t.Errorf("resolved AgentID = %q, want agent-1", got)
+ }
+}
diff --git a/internal/certificate/check_result_test.go b/internal/certificate/check_result_test.go
index 4637c77d..93a86c41 100644
--- a/internal/certificate/check_result_test.go
+++ b/internal/certificate/check_result_test.go
@@ -146,7 +146,13 @@ func TestProcessCheckResult_EachAlertRecoversWhenItsConditionClears(t *testing.T
bad := healthyScan()
tc.broken(bad)
svc.processCheckResult(ctx, monitor, bad, true)
- assert.Contains(t, alertTypes(drain()), tc.alertType)
+ fired := drain()
+ assert.Contains(t, alertTypes(fired), tc.alertType)
+ for _, e := range fired {
+ if e.eventType == event.CertificateAlert {
+ assert.Equal(t, monitor.AgentID, e.data["agent_id"], "the alert names the agent that scans the certificate")
+ }
+ }
svc.processCheckResult(ctx, monitor, bad, true)
assert.Empty(t, recoveredTypes(drain()), "a condition still present must not recover")
diff --git a/internal/certificate/service.go b/internal/certificate/service.go
index 16ff8d2a..553a3832 100644
--- a/internal/certificate/service.go
+++ b/internal/certificate/service.go
@@ -818,9 +818,9 @@ func (s *Service) evaluateAlerts(monitor *CertMonitor, result, previous *CertChe
return
}
- // withSNI tags alert payloads with the SNI so multi-vhost monitors on the
- // same host:port are distinguishable in notifications.
- withSNI := func(data map[string]interface{}) map[string]interface{} {
+ // withMonitor tags an alert payload with the scanning agent and the SNI that tells multi-vhost monitors apart.
+ withMonitor := func(data map[string]interface{}) map[string]interface{} {
+ data["agent_id"] = uid.Agent(monitor.AgentID)
if monitor.ServerName != "" {
data["server_name"] = monitor.ServerName
}
@@ -828,7 +828,7 @@ func (s *Service) evaluateAlerts(monitor *CertMonitor, result, previous *CertChe
}
recovery := func(alertType string) {
- data := withSNI(map[string]interface{}{
+ data := withMonitor(map[string]interface{}{
"monitor_id": monitor.ID,
"hostname": monitor.Hostname,
"port": monitor.Port,
@@ -855,7 +855,7 @@ func (s *Service) evaluateAlerts(monitor *CertMonitor, result, previous *CertChe
// Check chain validation alerts
if result.ChainValid != nil && !*result.ChainValid {
- s.emit(event.CertificateAlert, withSNI(map[string]interface{}{
+ s.emit(event.CertificateAlert, withMonitor(map[string]interface{}{
"monitor_id": monitor.ID,
"hostname": monitor.Hostname,
"port": monitor.Port,
@@ -868,7 +868,7 @@ func (s *Service) evaluateAlerts(monitor *CertMonitor, result, previous *CertChe
// Check hostname mismatch alerts
if result.HostnameMatch != nil && !*result.HostnameMatch {
- s.emit(event.CertificateAlert, withSNI(map[string]interface{}{
+ s.emit(event.CertificateAlert, withMonitor(map[string]interface{}{
"monitor_id": monitor.ID,
"hostname": monitor.Hostname,
"port": monitor.Port,
@@ -881,7 +881,7 @@ func (s *Service) evaluateAlerts(monitor *CertMonitor, result, previous *CertChe
// Only "revoked" triggers an alert. "unknown" and "error" are inconclusive
// and do not warrant a critical page.
if result.OCSPStatus == "revoked" {
- alertData := withSNI(map[string]interface{}{
+ alertData := withMonitor(map[string]interface{}{
"monitor_id": monitor.ID,
"hostname": monitor.Hostname,
"port": monitor.Port,
@@ -903,7 +903,7 @@ func (s *Service) evaluateAlerts(monitor *CertMonitor, result, previous *CertChe
// Check if certificate has expired
if result.NotAfter.Before(time.Now()) {
- s.emit(event.CertificateAlert, withSNI(map[string]interface{}{
+ s.emit(event.CertificateAlert, withMonitor(map[string]interface{}{
"monitor_id": monitor.ID,
"hostname": monitor.Hostname,
"port": monitor.Port,
@@ -944,7 +944,7 @@ func (s *Service) evaluateAlerts(monitor *CertMonitor, result, previous *CertChe
// Check if we need to escalate (fire alert at a new, lower threshold)
if monitor.LastAlertedThreshold == nil || *crossedThreshold < *monitor.LastAlertedThreshold {
- s.emit(event.CertificateAlert, withSNI(map[string]interface{}{
+ s.emit(event.CertificateAlert, withMonitor(map[string]interface{}{
"monitor_id": monitor.ID,
"hostname": monitor.Hostname,
"port": monitor.Port,
diff --git a/internal/eol/service.go b/internal/eol/service.go
index 00271b9f..524f0863 100644
--- a/internal/eol/service.go
+++ b/internal/eol/service.go
@@ -397,6 +397,7 @@ func (s *Service) forget(agentID string, now time.Time) {
EntityType: "agent",
EntityID: agentID,
EntityName: agentID,
+ AgentID: agentID,
Timestamp: now,
})
}
@@ -415,6 +416,7 @@ func hostEvent(a agent.Agent, identity Identity, support Support, now time.Time)
EntityType: "agent",
EntityID: a.AgentID,
EntityName: name,
+ AgentID: a.AgentID,
Message: supportMessage(identity, support),
Details: supportDetails(identity, support),
Timestamp: now,
diff --git a/internal/eol/service_test.go b/internal/eol/service_test.go
index b101ed81..aa4a6ca7 100644
--- a/internal/eol/service_test.go
+++ b/internal/eol/service_test.go
@@ -110,7 +110,7 @@ func TestEvaluateAllDebianJourney(t *testing.T) {
t.Fatalf("warning = %+v", warning)
}
if warning.Source != alert.SourceHost || warning.AlertType != AlertType ||
- warning.EntityType != "agent" || warning.EntityID != "web-1" || warning.EntityName != "web-1" {
+ warning.EntityType != "agent" || warning.EntityID != "web-1" || warning.EntityName != "web-1" || warning.AgentID != "web-1" {
t.Fatalf("event key = %+v", warning)
}
if want := "Debian GNU/Linux 11 (bullseye): free security support ends on 2026-08-31 (in 30 days)"; warning.Message != want {
@@ -263,7 +263,7 @@ func TestEvaluateAllResolvesAgentsThatLeft(t *testing.T) {
if len(got) != 1 {
t.Fatalf("events = %d, want 1", len(got))
}
- if !got[0].IsRecover || got[0].EntityID != "web-1" {
+ if !got[0].IsRecover || got[0].EntityID != "web-1" || got[0].AgentID != "web-1" {
t.Fatalf("event = %+v", got[0])
}
diff --git a/internal/resource/service.go b/internal/resource/service.go
index 166d7c81..9753afab 100644
--- a/internal/resource/service.go
+++ b/internal/resource/service.go
@@ -398,6 +398,7 @@ func (s *Service) evaluateAlerts(ctx context.Context, snap *ResourceSnapshot) {
"current_value": m.value,
"threshold": m.thresh,
"timestamp": now,
+ "agent_id": snap.AgentID,
})
case m.was && !m.is:
s.emit(event.ResourceRecovery, map[string]interface{}{
@@ -407,6 +408,7 @@ func (s *Service) evaluateAlerts(ctx context.Context, snap *ResourceSnapshot) {
"current_value": m.value,
"threshold": m.thresh,
"timestamp": now,
+ "agent_id": snap.AgentID,
})
}
}
diff --git a/internal/resource/service_test.go b/internal/resource/service_test.go
index f17cbb08..7959ff40 100644
--- a/internal/resource/service_test.go
+++ b/internal/resource/service_test.go
@@ -396,7 +396,9 @@ func TestService_evaluateAlerts_AlertEventFiredOnCPUTransition(t *testing.T) {
// No alert yet — only one breach.
assert.Empty(t, events)
- svc.evaluateAlerts(ctx, snap(containerID, 90.0, 0, 0))
+ remote := snap(containerID, 90.0, 0, 0)
+ remote.AgentID = "agent-1"
+ svc.evaluateAlerts(ctx, remote)
// Second breach triggers alert.
require.Len(t, events, 1)
assert.Equal(t, event.ResourceAlert, events[0].eventType)
@@ -404,6 +406,7 @@ func TestService_evaluateAlerts_AlertEventFiredOnCPUTransition(t *testing.T) {
assert.InDelta(t, 90.0, events[0].data["current_value"], 0.001)
assert.InDelta(t, 80.0, events[0].data["threshold"], 0.001)
assert.Equal(t, containerID, events[0].data["container_id"])
+ assert.Equal(t, "agent-1", events[0].data["agent_id"])
}
func TestService_evaluateAlerts_AlertEventFiredOnMemoryTransition(t *testing.T) {
diff --git a/internal/store/alerts.go b/internal/store/alerts.go
index 00e7a72a..28a2e698 100644
--- a/internal/store/alerts.go
+++ b/internal/store/alerts.go
@@ -31,21 +31,23 @@ func NewAlertStore(d *DB) *AlertStoreImpl {
const alertColumns = `id, source, alert_type, severity, status, message,
entity_type, entity_id, entity_name, details,
resolved_by_id, fired_at, resolved_at,
- acknowledged_at, acknowledged_by, escalated_at, created_at`
+ acknowledged_at, acknowledged_by, escalated_at, created_at, agent_id`
func (s *AlertStoreImpl) InsertAlert(ctx context.Context, a *alert.Alert) (string, error) {
a.ID = uid.New()
+ a.AgentID = uid.Agent(a.AgentID)
if a.CreatedAt.IsZero() {
a.CreatedAt = time.Now().UTC()
}
_, err := s.writer.Exec(ctx,
`INSERT INTO alerts (id, source, alert_type, severity, status, message,
entity_type, entity_id, entity_name, details,
- resolved_by_id, fired_at, resolved_at, created_at)
- VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
+ resolved_by_id, fired_at, resolved_at, created_at, agent_id)
+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
a.ID, a.Source, a.AlertType, a.Severity, a.Status, a.Message,
a.EntityType, a.EntityID, a.EntityName, a.Details,
nullableStrPtr(a.ResolvedByID), a.FiredAt.Unix(), nullableTime(a.ResolvedAt), a.CreatedAt.Unix(),
+ a.AgentID,
)
if err != nil {
return "", fmt.Errorf("insert alert: %w", err)
@@ -195,7 +197,7 @@ func scanAlertFromRow(scanner rowScanner) (*alert.Alert, error) {
&a.ID, &a.Source, &a.AlertType, &a.Severity, &a.Status, &a.Message,
&a.EntityType, &a.EntityID, &a.EntityName, &details,
&resolvedByID, &firedAt, &resolvedAt,
- &acknowledgedAt, &acknowledgedBy, &escalatedAt, &createdAt,
+ &acknowledgedAt, &acknowledgedBy, &escalatedAt, &createdAt, &a.AgentID,
)
if err != nil {
return nil, err
diff --git a/internal/store/alerts_test.go b/internal/store/alerts_test.go
index 5ad10b3d..e9404b9d 100644
--- a/internal/store/alerts_test.go
+++ b/internal/store/alerts_test.go
@@ -13,6 +13,7 @@ import (
"github.com/stretchr/testify/require"
"github.com/kolapsis/maintenant/internal/alert"
+ "github.com/kolapsis/maintenant/internal/uid"
)
func seedAlert(t *testing.T, s *AlertStoreImpl, entityID, status string, firedAt time.Time, resolvedAt *time.Time) string {
@@ -139,3 +140,34 @@ func TestAlertStore_ListUnacknowledgedActiveAlerts(t *testing.T) {
assert.Equal(t, open, alerts[0].ID, "an escalated alert is still unacknowledged")
require.NotNil(t, alerts[0].EscalatedAt)
}
+
+func TestAlertStore_AgentIDRoundTrips(t *testing.T) {
+ s := NewAlertStore(openTestDB(t))
+ ctx := context.Background()
+ now := time.Now().UTC()
+ const remote = "0190a1b2-c3d4-7e5f-8a9b-0c1d2e3f4a5b"
+
+ id, err := s.InsertAlert(ctx, &alert.Alert{
+ Source: alert.SourceContainer, AlertType: "health_unhealthy", Severity: alert.SeverityWarning,
+ Status: alert.StatusActive, Message: "unhealthy", EntityType: "container", EntityID: "c1", EntityName: "web",
+ FiredAt: now, AgentID: remote,
+ })
+ require.NoError(t, err)
+ local := seedAlert(t, s, "c2", alert.StatusActive, now, nil)
+
+ got, err := s.GetAlert(ctx, id)
+ require.NoError(t, err)
+ assert.Equal(t, remote, got.AgentID)
+
+ got, err = s.GetAlert(ctx, local)
+ require.NoError(t, err)
+ assert.Equal(t, uid.LocalAgent, got.AgentID, "an alert without an agent belongs to the local runtime")
+
+ active, err := s.ListActiveAlerts(ctx)
+ require.NoError(t, err)
+ byID := map[string]string{}
+ for _, a := range active {
+ byID[a.ID] = a.AgentID
+ }
+ assert.Equal(t, map[string]string{id: remote, local: uid.LocalAgent}, byID)
+}
diff --git a/internal/store/migration40_test.go b/internal/store/migration40_test.go
new file mode 100644
index 00000000..0607cb37
--- /dev/null
+++ b/internal/store/migration40_test.go
@@ -0,0 +1,40 @@
+// Copyright 2026 Benjamin Touchard (Kolapsis)
+// SPDX-License-Identifier: Apache-2.0
+
+package store
+
+import (
+ "context"
+ "io/fs"
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+
+ "github.com/kolapsis/maintenant/internal/uid"
+)
+
+func TestMigration40_ExistingAlertsBelongToTheLocalAgent(t *testing.T) {
+ db := openTestDB(t)
+ ctx := context.Background()
+
+ apply := func(direction string) {
+ t.Helper()
+ sqlText, err := fs.ReadFile(migrationFS,
+ "migrations/"+db.dialect.String()+"/40_alert_agent_id."+direction+".sql")
+ require.NoError(t, err)
+ _, err = db.ReadDB().ExecContext(ctx, string(sqlText))
+ require.NoError(t, err, direction)
+ }
+
+ apply("down")
+ _, err := db.ReadDB().ExecContext(ctx, db.dialect.Rebind(
+ `INSERT INTO alerts (id, source, alert_type, message, entity_type, entity_id, entity_name, fired_at)
+ VALUES (?, 'container', 'restart_loop', 'm', 'container', 'c1', 'web', 1)`), uid.New())
+ require.NoError(t, err)
+ apply("up")
+
+ var agentID string
+ require.NoError(t, db.ReadDB().QueryRowContext(ctx, `SELECT agent_id FROM alerts`).Scan(&agentID))
+ assert.Equal(t, uid.LocalAgent, agentID)
+}
diff --git a/internal/store/migrations/postgres/40_alert_agent_id.down.sql b/internal/store/migrations/postgres/40_alert_agent_id.down.sql
new file mode 100644
index 00000000..73808850
--- /dev/null
+++ b/internal/store/migrations/postgres/40_alert_agent_id.down.sql
@@ -0,0 +1 @@
+ALTER TABLE alerts DROP COLUMN agent_id;
diff --git a/internal/store/migrations/postgres/40_alert_agent_id.up.sql b/internal/store/migrations/postgres/40_alert_agent_id.up.sql
new file mode 100644
index 00000000..6734db6b
--- /dev/null
+++ b/internal/store/migrations/postgres/40_alert_agent_id.up.sql
@@ -0,0 +1 @@
+ALTER TABLE alerts ADD COLUMN agent_id TEXT NOT NULL DEFAULT '00000000-0000-0000-0000-000000000000';
diff --git a/internal/store/migrations/sqlite/40_alert_agent_id.down.sql b/internal/store/migrations/sqlite/40_alert_agent_id.down.sql
new file mode 100644
index 00000000..73808850
--- /dev/null
+++ b/internal/store/migrations/sqlite/40_alert_agent_id.down.sql
@@ -0,0 +1 @@
+ALTER TABLE alerts DROP COLUMN agent_id;
diff --git a/internal/store/migrations/sqlite/40_alert_agent_id.up.sql b/internal/store/migrations/sqlite/40_alert_agent_id.up.sql
new file mode 100644
index 00000000..6734db6b
--- /dev/null
+++ b/internal/store/migrations/sqlite/40_alert_agent_id.up.sql
@@ -0,0 +1 @@
+ALTER TABLE alerts ADD COLUMN agent_id TEXT NOT NULL DEFAULT '00000000-0000-0000-0000-000000000000';
diff --git a/internal/store/uuid_schema.sql b/internal/store/uuid_schema.sql
index 8fe37a5f..53c5bd2b 100644
--- a/internal/store/uuid_schema.sql
+++ b/internal/store/uuid_schema.sql
@@ -399,7 +399,8 @@ CREATE TABLE alerts (
created_at BIGINT NOT NULL DEFAULT 0,
acknowledged_at BIGINT,
acknowledged_by TEXT,
- escalated_at BIGINT
+ escalated_at BIGINT,
+ agent_id TEXT NOT NULL DEFAULT '00000000-0000-0000-0000-000000000000'
);
CREATE INDEX idx_alerts_status ON alerts(status);
CREATE INDEX idx_alerts_source_severity ON alerts(source, severity);
diff --git a/internal/update/container_adapter.go b/internal/update/container_adapter.go
index 690d866d..9f927812 100644
--- a/internal/update/container_adapter.go
+++ b/internal/update/container_adapter.go
@@ -116,6 +116,7 @@ func (a *ContainerServiceAdapter) detailsFor(c *container.Container, local map[s
func newContainerInfo(c *container.Container, d RuntimeDetails) ContainerInfo {
return ContainerInfo{
UID: c.ID,
+ AgentID: uid.Agent(c.AgentID),
ExternalID: c.ExternalID,
Name: c.Name,
Image: c.Image,
diff --git a/internal/update/container_adapter_test.go b/internal/update/container_adapter_test.go
index ca44beb0..2dd8dab2 100644
--- a/internal/update/container_adapter_test.go
+++ b/internal/update/container_adapter_test.go
@@ -144,6 +144,7 @@ func TestContainerServiceAdapter_AgentContainersUseWhatTheAgentReported(t *testi
byID := infosByID(t, adapter)
require.Contains(t, byID, "reported")
+ assert.Equal(t, agentID, byID["reported"].AgentID)
assert.Equal(t, map[string]string{"maintenant.update.enabled": "false"}, byID["reported"].Labels)
assert.Equal(t, []string{"nginx@sha256:running"}, byID["reported"].RepoDigests)
assert.NotContains(t, byID, "silent", "an agent container is scanned once its agent has reported its labels")
diff --git a/internal/update/scanner.go b/internal/update/scanner.go
index a0323470..79d58bfb 100644
--- a/internal/update/scanner.go
+++ b/internal/update/scanner.go
@@ -16,6 +16,7 @@ import (
// ContainerInfo holds the minimal container data needed for scanning.
type ContainerInfo struct {
UID string // canonical store PK (uid.Container) — used as alert EntityID
+ AgentID string
ExternalID string
Name string
Image string
diff --git a/internal/update/service.go b/internal/update/service.go
index b0657a67..fc5e9ce3 100644
--- a/internal/update/service.go
+++ b/internal/update/service.go
@@ -464,6 +464,7 @@ func (s *Service) runScan(ctx context.Context) {
"container_id": r.ContainerID,
"container_uid": containerByID[r.ContainerID].UID,
"container_name": r.ContainerName,
+ "agent_id": containerByID[r.ContainerID].AgentID,
"image": r.Image,
"current_tag": r.CurrentTag,
"latest_tag": r.LatestTag,
@@ -536,6 +537,7 @@ func (s *Service) runScan(ctx context.Context) {
"container_id": su.ContainerID,
"container_uid": containerUID,
"container_name": su.ContainerName,
+ "agent_id": containerByID[su.ContainerID].AgentID,
})
}
diff --git a/internal/update/service_commands_test.go b/internal/update/service_commands_test.go
index 7c92bae4..b066412a 100644
--- a/internal/update/service_commands_test.go
+++ b/internal/update/service_commands_test.go
@@ -245,7 +245,7 @@ func TestRunScan_PersistsThePreviousImageAndCarriesAlertOn(t *testing.T) {
Store: store,
Scanner: newTestScanner(reg, store),
Containers: stubLister{containers: []ContainerInfo{{
- UID: "uid1", ExternalID: "ctr1", Name: "web", Image: "nginx:1.24.0",
+ UID: "uid1", AgentID: "agent-1", ExternalID: "ctr1", Name: "web", Image: "nginx:1.24.0",
RepoDigests: []string{"nginx@sha256:running"},
Labels: map[string]string{"maintenant.update.alert_on": "critical"},
}}},
@@ -263,6 +263,7 @@ func TestRunScan_PersistsThePreviousImageAndCarriesAlertOn(t *testing.T) {
assert.Equal(t, "sha256:running", store.inserted[0].PreviousDigest)
require.NotNil(t, detected)
assert.Equal(t, AlertOnCritical, detected["alert_on"])
+ assert.Equal(t, "agent-1", detected["agent_id"])
assert.Equal(t, "docker pull nginx@sha256:running\n"+
"docker stop web && docker rm web\n"+
"docker run -d --name web nginx@sha256:running", detected["rollback_command"])
diff --git a/mkdocs.yml b/mkdocs.yml
index 4f87f42a..dc8be666 100644
--- a/mkdocs.yml
+++ b/mkdocs.yml
@@ -1,6 +1,6 @@
site_name: maintenant
site_description: The all-in-one monitoring dashboard your self-hosted stack deserves.
-site_url: https://kolapsis.github.io/maintenant/
+site_url: https://docs.maintenant.dev/
repo_name: kolapsis/maintenant
repo_url: https://github.com/kolapsis/maintenant
@@ -52,7 +52,11 @@ markdown_extensions:
- toc:
permalink: true
- pymdownx.details
- - pymdownx.superfences
+ - pymdownx.superfences:
+ custom_fences:
+ - name: mermaid
+ class: mermaid
+ format: !!python/name:pymdownx.superfences.fence_code_format
- pymdownx.highlight:
anchor_linenums: true
line_spans: __span
@@ -96,6 +100,7 @@ nav:
- Vultr Deployment: guides/vultr.md
- Docker Swarm Deployment: guides/swarm.md
- PostgreSQL Storage: guides/postgresql.md
+ - AI Agents: guides/ai-agents.md
- API:
- API Reference: api/reference.md
- Security: security.md
From 9a29b25c282fce63db1deb8d6fed76ce4895dc3b Mon Sep 17 00:00:00 2001
From: Benjamin
Date: Thu, 1 Oct 2026 21:37:40 +0200
Subject: [PATCH 2/3] test(store): undo migration 40 in the PostgreSQL
concurrent catch-up test
---
internal/store/migrations_postgres_test.go | 1 +
1 file changed, 1 insertion(+)
diff --git a/internal/store/migrations_postgres_test.go b/internal/store/migrations_postgres_test.go
index 42925d98..b343876c 100644
--- a/internal/store/migrations_postgres_test.go
+++ b/internal/store/migrations_postgres_test.go
@@ -124,6 +124,7 @@ func TestMigratePostgres_ConcurrentCatchUp(t *testing.T) {
"ALTER TABLE heartbeats ADD COLUMN active INTEGER NOT NULL DEFAULT 1", // 39
"CREATE INDEX idx_heartbeat_status_deadline ON heartbeats(status, next_deadline_at) WHERE active=1", // 39
"CREATE INDEX idx_heartbeat_active ON heartbeats(active)", // 39
+ "ALTER TABLE alerts DROP COLUMN agent_id", // 40
} {
_, err = db.ReadDB().Exec(undo)
require.NoError(t, err, undo)
From b6a09cc6c85d03bcf56f5bdaf13f0940abe41608 Mon Sep 17 00:00:00 2001
From: Benjamin
Date: Thu, 1 Oct 2026 21:37:40 +0200
Subject: [PATCH 3/3] fix(agent): close the push stream under the send lock
CloseSend raced with a snapshot still being sent by the spool when the
stream ended, which gRPC does not allow.
---
internal/agent/client.go | 2 ++
1 file changed, 2 insertions(+)
diff --git a/internal/agent/client.go b/internal/agent/client.go
index fbe4e2f4..2dd7a33c 100644
--- a/internal/agent/client.go
+++ b/internal/agent/client.go
@@ -157,6 +157,8 @@ func (ps *PushStream) SendResult(res *agentpb.CommandResult) error {
// Close signals the end of the send side of the stream.
func (ps *PushStream) Close() {
+ ps.mu.Lock()
+ defer ps.mu.Unlock()
_ = ps.stream.CloseSend()
}
|