List alert sources with their status and last arrival on Team -> Sources

Like Team -> Switches, the page is now a table: a status badge (Active
if the key posted within a day, Quiet if it has but not lately, Never
used), when it last posted a webhook, when an alert last arrived on it,
how many distinct alerts it refreshed in the last 24 hours, and when it
was created. Adding a source moved into a "New source" sheet, and owners
can rename one from its row.

"Last alert" and the count needed alerts to remember which source they
came in on, which they never did, so migration 010 adds
alerts.integration_id and every accepted payload stamps it. Last sender
wins when two sources post the same fingerprint. It is not backfilled: a
NULL says "before this was recorded" rather than guessing, and it heals
by itself as Alertmanager re-sends each alert every repeat_interval.
Revoking a source keeps its alerts, unattributed.

Last webhook and last alert are separate on purpose: a payload with
nothing usable in it stamps the first and not the second. The Quiet
threshold is a fixed day, a colour and not an alarm, since silence that
should page is what dead man's switches are for.

The counts are indexed subqueries (alerts_integration_idx) rather than a
join, which would read every alert a source ever delivered.

API: the integrations list gains status, last_alert_at and alerts_24h,
and PATCH /api/teams/{id}/integrations/{id} renames. Both are additive;
terdut-tui needs nothing.
This commit is contained in:
Niklas Ye
2026-09-26 07:52:07 +02:00
parent e8d45f9d3d
commit d675f8ec9b
9 changed files with 439 additions and 69 deletions
+15 -10
View File
@@ -78,7 +78,7 @@ type ingested struct {
// post, and which team the alerts belong to.
func handleIntegrationWebhook(db *sql.DB, notify NotifyConfig) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
teamID, err := teamIDForKey(r.Context(), db, chi.URLParam(r, "key"))
src, err := sourceForKey(r.Context(), db, chi.URLParam(r, "key"))
if err != nil {
if errors.Is(err, errUnknownIntegration) {
// 401 and not 404: the path is real, the key is not, and a
@@ -90,11 +90,12 @@ func handleIntegrationWebhook(db *sql.DB, notify NotifyConfig) http.HandlerFunc
respond(w, http.StatusInternalServerError, errResp("internal error"))
return
}
receiveWebhook(w, r, db, notify, teamID)
receiveWebhook(w, r, db, notify, src)
}
}
func receiveWebhook(w http.ResponseWriter, r *http.Request, db *sql.DB, notify NotifyConfig, teamID int64) {
func receiveWebhook(w http.ResponseWriter, r *http.Request, db *sql.DB, notify NotifyConfig, src alertSource) {
teamID := src.teamID
var payload amPayload
if err := decodeJSON(r, &payload); err != nil {
respond(w, http.StatusBadRequest, errResp("invalid payload"))
@@ -104,7 +105,7 @@ func receiveWebhook(w http.ResponseWriter, r *http.Request, db *sql.DB, notify N
// Alertmanager retries anything that is not 2xx, and a retry of a payload
// we failed to store is more useful than an error it cannot act on — so
// failures are logged, not surfaced.
if err := ingest(r.Context(), db, notify, teamID, payload); err != nil {
if err := ingest(r.Context(), db, notify, src, payload); err != nil {
log.Printf("webhook ingest (team %d, group %q): %v", teamID, payload.GroupKey, err)
}
@@ -114,7 +115,8 @@ func receiveWebhook(w http.ResponseWriter, r *http.Request, db *sql.DB, notify N
// ingest stores a payload's alerts and reconciles the incident for its group.
// The whole payload is one transaction: an incident that opened but whose alerts
// failed to link would be a work item nobody could act on.
func ingest(ctx context.Context, db *sql.DB, notify NotifyConfig, teamID int64, payload amPayload) error {
func ingest(ctx context.Context, db *sql.DB, notify NotifyConfig, src alertSource, payload amPayload) error {
teamID := src.teamID
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
@@ -129,7 +131,7 @@ func ingest(ctx context.Context, db *sql.DB, notify NotifyConfig, teamID int64,
return err
}
accepted, err := upsertAlerts(ctx, tx, deadman, teamID, payload.Alerts)
accepted, err := upsertAlerts(ctx, tx, deadman, src, payload.Alerts)
if err != nil {
return err
}
@@ -186,7 +188,8 @@ func ingest(ctx context.Context, db *sql.DB, notify NotifyConfig, teamID int64,
// upsertAlerts stores each alert of a payload and reports what changed. Payloads
// the ordering guard rejected are left out entirely.
func upsertAlerts(ctx context.Context, tx *sql.Tx, deadman deadmanSet, teamID int64, alerts []amAlert) ([]ingested, error) {
func upsertAlerts(ctx context.Context, tx *sql.Tx, deadman deadmanSet, src alertSource, alerts []amAlert) ([]ingested, error) {
teamID := src.teamID
now := time.Now().Unix()
accepted := make([]ingested, 0, len(alerts))
@@ -242,8 +245,8 @@ func upsertAlerts(ctx context.Context, tx *sql.Tx, deadman deadmanSet, teamID in
if _, err := tx.ExecContext(ctx, `
INSERT INTO alerts
(team_id, fingerprint, name, status, labels, annotations, starts_at, ends_at,
generator_url, received_at, resolution_source)
VALUES ($1, $2, $3, $4, $5::jsonb, $6::jsonb, $7, $8, $9, $10, $11)
generator_url, received_at, resolution_source, integration_id)
VALUES ($1, $2, $3, $4, $5::jsonb, $6::jsonb, $7, $8, $9, $10, $11, $12)
ON CONFLICT (team_id, fingerprint) DO UPDATE SET
status = excluded.status,
labels = excluded.labels,
@@ -257,6 +260,8 @@ func upsertAlerts(ctx context.Context, tx *sql.Tx, deadman deadmanSet, teamID in
-- is a breaking API change — see models.Alert.ReceivedAt.
received_at = excluded.received_at,
resolution_source = excluded.resolution_source,
-- Last sender wins; see migration 010.
integration_id = excluded.integration_id,
-- A re-fire makes the alert current again, so it leaves the archive.
archived_at = CASE WHEN excluded.status = 'firing'
THEN NULL ELSE alerts.archived_at END
@@ -267,7 +272,7 @@ func upsertAlerts(ctx context.Context, tx *sql.Tx, deadman deadmanSet, teamID in
teamID, a.Fingerprint, name, a.Status,
string(labelsJSON), string(annotationsJSON),
a.StartsAt.Unix(), endsAtUnix,
a.GeneratorURL, now, resolutionSource,
a.GeneratorURL, now, resolutionSource, src.integrationID,
); err != nil {
return nil, err
}
+1
View File
@@ -151,6 +151,7 @@ func NewRouter(db *sql.DB, notify NotifyConfig, cfg config.Config) http.Handler
// Integrations: where a team's alerts come in, and the key that says so.
r.Get("/api/teams/{teamID}/integrations", handleListIntegrations(db))
r.Post("/api/teams/{teamID}/integrations", handleCreateIntegration(db, notify.PublicURL))
r.Patch("/api/teams/{teamID}/integrations/{integrationID}", handleRenameIntegration(db))
r.Delete("/api/teams/{teamID}/integrations/{integrationID}", handleDeleteIntegration(db))
// The rota is per team. /api/schedule/current is the exception: it
+171
View File
@@ -0,0 +1,171 @@
package api_test
import (
"bytes"
"net/http"
"strings"
"testing"
"time"
)
// listSources reads a team's alert sources as the Sources page does.
func listSources(t *testing.T, tm teamFixture) []map[string]any {
t.Helper()
return list(t, tm.call(http.MethodGet, "/api/teams/"+id64(tm.id)+"/integrations", nil))
}
// addSource mints a second source in a team and returns its key.
func addSource(t *testing.T, tm teamFixture, name string) string {
t.Helper()
var out struct {
Key string `json:"key"`
}
decode(t, tm.call(http.MethodPost, "/api/teams/"+id64(tm.id)+"/integrations",
map[string]string{"name": name}), &out)
return out.Key
}
// A source that has never posted is "never", with nothing to say about alerts.
func TestSources_NeverUsedIsBlank(t *testing.T) {
s := newTS(t)
tm := newTeam(t, s, "red")
got := listSources(t, tm)
if len(got) != 1 {
t.Fatalf("expected 1 source, got %d", len(got))
}
src := got[0]
if src["status"] != "never" || src["last_used_at"] != nil || src["last_alert_at"] != nil {
t.Errorf("a source nobody has posted on should be blank, got %v", src)
}
if src["alerts_24h"].(float64) != 0 {
t.Errorf("alerts_24h = %v, want 0", src["alerts_24h"])
}
}
// Each source is credited with what arrived on its own key, and only that.
func TestSources_AlertsAreAttributedToTheirSource(t *testing.T) {
s := newTS(t)
tm := newTeam(t, s, "red")
second := addSource(t, tm, "staging")
postToIntegration(t, s, tm.key, "fp-1", "DiskFull")
postToIntegration(t, s, tm.key, "fp-2", "CPUHot")
got := listSources(t, tm)
first, other := got[0], got[1]
if first["status"] != "active" || first["last_used_at"] == nil || first["last_alert_at"] == nil {
t.Errorf("the source that posted should be active with timestamps, got %v", first)
}
if first["alerts_24h"].(float64) != 2 {
t.Errorf("alerts_24h = %v, want 2", first["alerts_24h"])
}
if other["status"] != "never" || other["alerts_24h"].(float64) != 0 {
t.Errorf("the other source should be untouched, got %v", other)
}
// Re-sending the same alert on the other key moves it: last sender wins.
postToIntegration(t, s, second, "fp-1", "DiskFull")
got = listSources(t, tm)
if got[0]["alerts_24h"].(float64) != 1 || got[1]["alerts_24h"].(float64) != 1 {
t.Errorf("fp-1 should have moved to the second source, got %v and %v",
got[0]["alerts_24h"], got[1]["alerts_24h"])
}
}
// A payload with no alerts in it is a webhook, not an alert: the source was
// heard from, and nothing arrived.
func TestSources_EmptyPayloadStampsUseButNotAlert(t *testing.T) {
s := newTS(t)
tm := newTeam(t, s, "red")
resp, err := http.Post(s.URL+"/api/integrations/"+tm.key+"/alertmanager",
"application/json", bytes.NewReader([]byte(`{"version":"4","status":"firing","alerts":[]}`)))
if err != nil {
t.Fatalf("post: %v", err)
}
resp.Body.Close()
src := listSources(t, tm)[0]
if src["status"] != "active" || src["last_alert_at"] != nil {
t.Errorf("want active with no alert yet, got %v", src)
}
}
// Quiet is "has posted, not lately"; the alert counter forgets after a day but
// the last alert's timestamp is kept.
func TestSources_QuietAfterADay(t *testing.T) {
s := newTS(t)
tm := newTeam(t, s, "red")
postToIntegration(t, s, tm.key, "fp-1", "DiskFull")
old := time.Now().Add(-48 * time.Hour).Unix()
s.exec(t, "UPDATE integrations SET last_used_at = $1", old)
s.exec(t, "UPDATE alerts SET received_at = $1 WHERE fingerprint = 'fp-1'", old)
src := listSources(t, tm)[0]
if src["status"] != "quiet" {
t.Errorf("status = %v, want quiet", src["status"])
}
if src["alerts_24h"].(float64) != 0 {
t.Errorf("alerts_24h = %v, want 0", src["alerts_24h"])
}
if src["last_alert_at"] == nil {
t.Error("last_alert_at should survive the day")
}
}
// Revoking a source does not take its alerts with it.
func TestSources_RevokeKeepsTheAlerts(t *testing.T) {
s := newTS(t)
tm := newTeam(t, s, "red")
postToIntegration(t, s, tm.key, "fp-1", "DiskFull")
id := int64(listSources(t, tm)[0]["id"].(float64))
resp := tm.call(http.MethodDelete, "/api/teams/"+id64(tm.id)+"/integrations/"+id64(id), nil)
resp.Body.Close()
if resp.StatusCode != http.StatusNoContent {
t.Fatalf("revoke: %d", resp.StatusCode)
}
if got := len(list(t, tm.call(http.MethodGet, "/api/alerts", nil))); got != 1 {
t.Errorf("the alert should outlive its source, got %d alerts", got)
}
}
// Renaming is an owner's, scoped to the team, and does not touch the key.
func TestSources_Rename(t *testing.T) {
s := newTS(t)
tm := newTeam(t, s, "red")
other := newTeam(t, s, "blue")
id := int64(listSources(t, tm)[0]["id"].(float64))
path := "/api/teams/" + id64(tm.id) + "/integrations/" + id64(id)
resp := tm.call(http.MethodPatch, path, map[string]string{"name": " prod "})
resp.Body.Close()
if resp.StatusCode != http.StatusNoContent {
t.Fatalf("rename: %d", resp.StatusCode)
}
if name := listSources(t, tm)[0]["name"]; name != "prod" {
t.Errorf("name = %q, want it trimmed to prod", name)
}
postToIntegration(t, s, tm.key, "fp-1", "DiskFull") // the old key still works
for name, body := range map[string]map[string]string{
"empty": {"name": " "},
"too long": {"name": strings.Repeat("x", 101)},
} {
resp := tm.call(http.MethodPatch, path, body)
resp.Body.Close()
if resp.StatusCode != http.StatusBadRequest {
t.Errorf("%s name: expected 400, got %d", name, resp.StatusCode)
}
}
// Another team's owner cannot reach it.
resp = other.call(http.MethodPatch, "/api/teams/"+id64(other.id)+"/integrations/"+id64(id),
map[string]string{"name": "mine now"})
resp.Body.Close()
if resp.StatusCode != http.StatusNotFound {
t.Errorf("renaming another team's source: expected 404, got %d", resp.StatusCode)
}
}
+126 -19
View File
@@ -355,8 +355,42 @@ func isLastTeamOwner(ctx context.Context, db *sql.DB, teamID, userID int64) (boo
// Integrations
// ---------------------------------------------------------------------------
// handleListIntegrations lists a team's integrations. Never the keys: those
// exist in plaintext only in the response that created them.
// sourceQuietAfter is how long a source may go without posting before the
// Sources page calls it quiet rather than active. A day is longer than any
// repeat_interval worth having, so an Alertmanager that is up and has anything
// firing never crosses it; a source with nothing firing may, and that is a
// reason to look, not proof of a fault — which is why this is a colour and not
// an alarm. Dead man's switches are where silence pages.
const sourceQuietAfter = 24 * time.Hour
const (
sourceActive = "active"
sourceQuiet = "quiet"
sourceNever = "never"
)
// integrationStatus is an integration as the Sources page shows it.
type integrationStatus struct {
models.Integration
// Status is active when the key posted within sourceQuietAfter, quiet when
// it has posted but not lately, never when it has not posted at all.
Status string `json:"status"`
// LastAlertAt is when an alert last arrived on this source, which is not the
// same as when it last posted: a payload with nothing usable in it stamps
// last_used_at and not this. Absent until an alert has arrived since
// migration 010 started recording it.
LastAlertAt *time.Time `json:"last_alert_at,omitempty"`
// Alerts24h counts the distinct alerts this source refreshed in the last
// day. An alert re-sent every few hours counts once, not once per re-send.
Alerts24h int64 `json:"alerts_24h"`
}
// handleListIntegrations lists a team's integrations with what each has been
// delivering. Never the keys: those exist in plaintext only in the response that
// created them.
func handleListIntegrations(db *sql.DB) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
teamID, ok := teamParam(w, r)
@@ -367,28 +401,45 @@ func handleListIntegrations(db *sql.DB) http.HandlerFunc {
return
}
now := time.Now()
rows, err := db.QueryContext(r.Context(), `
SELECT id, team_id, kind, name, created_at, last_used_at
FROM integrations
WHERE team_id = $1
ORDER BY id`, teamID)
SELECT i.id, i.team_id, i.kind, i.name, i.created_at, i.last_used_at,
-- Scalar subqueries, not a join and GROUP BY: each is a
-- single range over alerts_integration_idx, where the join
-- would read every alert a source ever delivered.
(SELECT MAX(received_at) FROM alerts WHERE integration_id = i.id),
(SELECT COUNT(*) FROM alerts
WHERE integration_id = i.id AND received_at >= $2)
FROM integrations i
WHERE i.team_id = $1
ORDER BY i.id`, teamID, now.Add(-sourceQuietAfter).Unix())
if err != nil {
respond(w, http.StatusInternalServerError, errResp("internal error"))
return
}
defer rows.Close()
integrations := []models.Integration{}
integrations := []integrationStatus{}
for rows.Next() {
var i models.Integration
var i integrationStatus
var created int64
var lastUsed *int64
if err := rows.Scan(&i.ID, &i.TeamID, &i.Kind, &i.Name, &created, &lastUsed); err != nil {
var lastUsed, lastAlert *int64
if err := rows.Scan(&i.ID, &i.TeamID, &i.Kind, &i.Name, &created, &lastUsed,
&lastAlert, &i.Alerts24h); err != nil {
respond(w, http.StatusInternalServerError, errResp("internal error"))
return
}
i.CreatedAt = time.Unix(created, 0).UTC()
i.LastUsedAt = unixPtr(lastUsed)
i.LastAlertAt = unixPtr(lastAlert)
switch {
case i.LastUsedAt == nil:
i.Status = sourceNever
case now.Sub(*i.LastUsedAt) > sourceQuietAfter:
i.Status = sourceQuiet
default:
i.Status = sourceActive
}
integrations = append(integrations, i)
}
if err := rows.Err(); err != nil {
@@ -457,6 +508,54 @@ func handleCreateIntegration(db *sql.DB, publicURL string) http.HandlerFunc {
}
}
// handleRenameIntegration renames a source. The key is untouched, so nothing
// posting with it notices.
func handleRenameIntegration(db *sql.DB) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
teamID, ok := teamParam(w, r)
if !ok {
return
}
if !requireTeamOwner(w, r, teamID) {
return
}
id, err := strconv.ParseInt(chi.URLParam(r, "integrationID"), 10, 64)
if err != nil {
respond(w, http.StatusBadRequest, errResp("invalid integration id"))
return
}
var req struct {
Name string `json:"name"`
}
if err := decodeJSON(r, &req); err != nil {
respond(w, http.StatusBadRequest, errResp("invalid request body"))
return
}
req.Name = strings.TrimSpace(req.Name)
if req.Name == "" {
respond(w, http.StatusBadRequest, errResp("name is required"))
return
}
if len(req.Name) > 100 {
respond(w, http.StatusBadRequest, errResp("name is too long"))
return
}
res, err := db.ExecContext(r.Context(),
"UPDATE integrations SET name = $1 WHERE id = $2 AND team_id = $3", req.Name, id, teamID)
if err != nil {
respond(w, http.StatusInternalServerError, errResp("internal error"))
return
}
if n, _ := res.RowsAffected(); n == 0 {
respond(w, http.StatusNotFound, errResp("not found"))
return
}
w.WriteHeader(http.StatusNoContent)
}
}
func handleDeleteIntegration(db *sql.DB) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
teamID, ok := teamParam(w, r)
@@ -493,24 +592,32 @@ func integrationPath(key, kind string) string {
return "/api/integrations/" + key + "/" + kind
}
// teamIDForKey resolves an integration key to its team, and stamps the key's
// alertSource is who an arriving webhook is from: the integration whose key it
// used, and the team that integration puts its alerts in.
type alertSource struct {
integrationID int64
teamID int64
}
// sourceForKey resolves an integration key to its source, and stamps the key's
// last use. An unknown key is not an error worth distinguishing: the caller is
// told nothing beyond "no".
func teamIDForKey(ctx context.Context, db *sql.DB, key string) (int64, error) {
var teamID int64
func sourceForKey(ctx context.Context, db *sql.DB, key string) (alertSource, error) {
var src alertSource
err := db.QueryRowContext(ctx,
"SELECT team_id FROM integrations WHERE key_hash = $1", hashToken(key)).Scan(&teamID)
"SELECT id, team_id FROM integrations WHERE key_hash = $1", hashToken(key)).
Scan(&src.integrationID, &src.teamID)
if errors.Is(err, sql.ErrNoRows) {
return 0, errUnknownIntegration
return alertSource{}, errUnknownIntegration
}
if err != nil {
return 0, err
return alertSource{}, err
}
// Best effort, like an API key's: a failed stamp must not reject an alert.
db.ExecContext(ctx, //nolint:errcheck
"UPDATE integrations SET last_used_at = $1 WHERE key_hash = $2",
time.Now().Unix(), hashToken(key))
return teamID, nil
"UPDATE integrations SET last_used_at = $1 WHERE id = $2",
time.Now().Unix(), src.integrationID)
return src, nil
}
var errUnknownIntegration = errors.New("unknown integration key")