diff --git a/cmd/relay/README.md b/cmd/relay/README.md index 9d92640a..acf72cef 100644 --- a/cmd/relay/README.md +++ b/cmd/relay/README.md @@ -116,6 +116,15 @@ Be sure to double-check bandwidth usage and pricing if running a public relay! B The relay admin interface has flexibility for many situations, but in some operational incidents it may be necessary to run SQL commands to do cleanups. This should be done when the relay is not actively operating. It is also recommended to run SQL commands in a transaction that can be rolled back in case of a typo or mistake. +Account-limit alerts can be silenced for individual PDS hosts through the admin API. This state is persisted on the host row and forwarded to sibling relays: + + curl -u admin:$RELAY_ADMIN_PASSWORD \ + -H 'Content-Type: application/json' \ + -d '{"host":"pds.example.com","account_limit_alerts_silenced":true}' \ + http://localhost:2470/admin/pds/changeLimits + +Set `account_limit_alerts_silenced` to `false` to resume alerting. + On the public web, you should probably run the relay behind a load-balancer or reverse proxy like `haproxy` or `caddy`, which manages TLS and can have various HTTP limits and behaviors configured. Remember that WebSocket support is required. The relay does not resolve atproto handles, but it does do DNS resolutions for hostnames, and may do a burst of resolutions at startup. Note that the go runtime may have an internal DNS implementation enabled (this is the default for the Dockerfile). The relay *will* do a large number of DID resolutions, particularly calls to the PLC directory, and particularly after a process restart when the in-process identity cache is warming up. diff --git a/cmd/relay/alerts.go b/cmd/relay/alerts.go index ece71b5f..ecb3375d 100644 --- a/cmd/relay/alerts.go +++ b/cmd/relay/alerts.go @@ -11,10 +11,10 @@ import ( "net/http" "net/url" "strings" - "sync" "time" "github.com/bluesky-social/indigo/cmd/relay/relay" + "github.com/bluesky-social/indigo/cmd/relay/relay/models" "github.com/bluesky-social/indigo/util/cliutil" "github.com/labstack/echo/v4" @@ -58,8 +58,9 @@ type Alerter interface { SendAlert(ctx context.Context, msg AlertMessage) error } -type accountLimitUsageLister interface { - ListHostsApproachingAccountLimit(ctx context.Context, threshold float64, limit int) ([]relay.HostAccountLimitUsage, error) +type accountLimitAlertClaimer interface { + ClaimDueAccountLimitAlerts(ctx context.Context, threshold float64, repeatInterval time.Duration, limit int) (*relay.HostAccountLimitAlertClaims, error) + RecordHostAccountLimitAlertSent(ctx context.Context, hostname string, state models.HostAccountLimitAlertState, sentAt time.Time) error } type AccountLimitAlertSentState struct { @@ -161,19 +162,16 @@ func (c AccountLimitAlertMonitorConfig) Validate() error { type AccountLimitAlertMonitor struct { logger *slog.Logger - lister accountLimitUsageLister + claimer accountLimitAlertClaimer alerter Alerter config AccountLimitAlertMonitorConfig - lk sync.Mutex - now func() time.Time - lastAlert map[uint64]time.Time - aboveThreshold map[uint64]relay.HostAccountLimitUsage + now func() time.Time } -func NewAccountLimitAlertMonitor(logger *slog.Logger, lister accountLimitUsageLister, alerter Alerter, config AccountLimitAlertMonitorConfig) (*AccountLimitAlertMonitor, error) { - if lister == nil { - return nil, fmt.Errorf("account limit alert lister is required") +func NewAccountLimitAlertMonitor(logger *slog.Logger, claimer accountLimitAlertClaimer, alerter Alerter, config AccountLimitAlertMonitorConfig) (*AccountLimitAlertMonitor, error) { + if claimer == nil { + return nil, fmt.Errorf("account limit alert claimer is required") } if alerter == nil { return nil, fmt.Errorf("account limit alerter is required") @@ -185,13 +183,11 @@ func NewAccountLimitAlertMonitor(logger *slog.Logger, lister accountLimitUsageLi return nil, err } return &AccountLimitAlertMonitor{ - logger: logger.With("system", "account-limit-alerts"), - lister: lister, - alerter: alerter, - config: config, - now: time.Now, - lastAlert: make(map[uint64]time.Time), - aboveThreshold: make(map[uint64]relay.HostAccountLimitUsage), + logger: logger.With("system", "account-limit-alerts"), + claimer: claimer, + alerter: alerter, + config: config, + now: time.Now, }, nil } @@ -249,103 +245,38 @@ func jitteredCheckInterval(interval, jitter time.Duration, randInt63n func(int64 } func (m *AccountLimitAlertMonitor) Check(ctx context.Context) error { - usages, err := m.lister.ListHostsApproachingAccountLimit(ctx, m.config.Threshold, 0) + claims, err := m.claimer.ClaimDueAccountLimitAlerts(ctx, m.config.Threshold, m.config.RepeatInterval, 0) if err != nil { return err } - now := m.now() - m.lk.Lock() - currentAbove := make(map[uint64]relay.HostAccountLimitUsage, len(usages)) - due := make([]relay.HostAccountLimitUsage, 0, len(usages)) - for _, usage := range usages { - currentAbove[usage.ID] = usage - last, alertedBefore := m.lastAlert[usage.ID] - _, wasAbove := m.aboveThreshold[usage.ID] - repeatDue := !alertedBefore || now.Sub(last) >= m.config.RepeatInterval - if !wasAbove || repeatDue { - due = append(due, usage) - } - } - - recovered := make([]relay.HostAccountLimitUsage, 0) - for id := range m.aboveThreshold { - if _, ok := currentAbove[id]; !ok { - recovered = append(recovered, m.aboveThreshold[id]) - } - } - - for id, usage := range currentAbove { - m.aboveThreshold[id] = usage - } - m.lk.Unlock() - - for _, usage := range recovered { + for _, usage := range claims.Recoveries { if err := m.alerter.SendAlert(ctx, formatAccountLimitRecoveryAlert(m.config, usage)); err != nil { return err } - if err := m.recordAccountLimitAlertSent(ctx, accountLimitAlertSentState(AccountLimitAlertKindRecovery, usage, now), true); err != nil { - return err - } + m.forwardAccountLimitAlertSent(ctx, accountLimitAlertSentState(AccountLimitAlertKindRecovery, usage, usage.AlertSentAt)) } - for start := 0; start < len(due); start += m.config.HostsPerMessage { + for start := 0; start < len(claims.Warnings); start += m.config.HostsPerMessage { end := start + m.config.HostsPerMessage - if end > len(due) { - end = len(due) + if end > len(claims.Warnings) { + end = len(claims.Warnings) } - batch := due[start:end] + batch := claims.Warnings[start:end] if err := m.alerter.SendAlert(ctx, formatAccountLimitAlert(m.config, batch)); err != nil { return err } for _, usage := range batch { - if err := m.recordAccountLimitAlertSent(ctx, accountLimitAlertSentState(AccountLimitAlertKindWarning, usage, now), true); err != nil { - return err - } + m.forwardAccountLimitAlertSent(ctx, accountLimitAlertSentState(AccountLimitAlertKindWarning, usage, usage.AlertSentAt)) } } return nil } -func (m *AccountLimitAlertMonitor) RecordAccountLimitAlertSent(state AccountLimitAlertSentState) error { - return m.recordAccountLimitAlertSent(context.Background(), state, false) -} - -func (m *AccountLimitAlertMonitor) recordAccountLimitAlertSent(ctx context.Context, state AccountLimitAlertSentState, notifySiblings bool) error { - now := m.now() - if err := state.normalize(now); err != nil { - return err - } - - usage := relay.HostAccountLimitUsage{ - ID: state.HostID, - Hostname: state.Hostname, - AccountCount: state.AccountCount, - AccountLimit: state.AccountLimit, - Usage: state.Usage, - } - - m.lk.Lock() - switch state.Kind { - case AccountLimitAlertKindWarning: - if last, ok := m.lastAlert[state.HostID]; !ok || state.SentAt.After(last) { - m.lastAlert[state.HostID] = state.SentAt - } - m.aboveThreshold[state.HostID] = usage - case AccountLimitAlertKindRecovery: - if last, ok := m.lastAlert[state.HostID]; ok && last.After(state.SentAt) { - m.lk.Unlock() - return nil - } - delete(m.lastAlert, state.HostID) - delete(m.aboveThreshold, state.HostID) - } - m.lk.Unlock() - - if notifySiblings && m.config.SentCallback != nil { +func (m *AccountLimitAlertMonitor) forwardAccountLimitAlertSent(ctx context.Context, state AccountLimitAlertSentState) { + if m.config.SentCallback != nil { m.config.SentCallback(ctx, state) } - return nil } func (s *Service) handleAdminRecordAccountLimitAlertSent(c echo.Context) error { @@ -353,10 +284,14 @@ func (s *Service) handleAdminRecordAccountLimitAlertSent(c echo.Context) error { if err := c.Bind(&body); err != nil { return echo.NewHTTPError(http.StatusBadRequest, fmt.Sprintf("invalid body: %s", err)) } - if s.accountLimitAlertRecorder == nil { - return c.JSON(http.StatusOK, map[string]any{"success": "true"}) + if err := body.normalize(time.Now().UTC()); err != nil { + return echo.NewHTTPError(http.StatusBadRequest, err.Error()) + } + state := models.HostAccountLimitAlertStateWarning + if body.Kind == AccountLimitAlertKindRecovery { + state = models.HostAccountLimitAlertStateOK } - if err := s.accountLimitAlertRecorder.RecordAccountLimitAlertSent(body); err != nil { + if err := s.relay.RecordHostAccountLimitAlertSent(c.Request().Context(), body.Hostname, state, body.SentAt); err != nil { return echo.NewHTTPError(http.StatusBadRequest, err.Error()) } return c.JSON(http.StatusOK, map[string]any{"success": "true"}) diff --git a/cmd/relay/alerts_test.go b/cmd/relay/alerts_test.go index 01462dfa..62fdbe45 100644 --- a/cmd/relay/alerts_test.go +++ b/cmd/relay/alerts_test.go @@ -11,19 +11,35 @@ import ( "time" "github.com/bluesky-social/indigo/cmd/relay/relay" + "github.com/bluesky-social/indigo/cmd/relay/relay/models" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) -type fakeAccountLimitUsageLister struct { - usages []relay.HostAccountLimitUsage +type fakeAccountLimitAlertClaimer struct { + claims []relay.HostAccountLimitAlertClaims lastLimit int + calls int + records []string } -func (f *fakeAccountLimitUsageLister) ListHostsApproachingAccountLimit(ctx context.Context, threshold float64, limit int) ([]relay.HostAccountLimitUsage, error) { +func (f *fakeAccountLimitAlertClaimer) ClaimDueAccountLimitAlerts(ctx context.Context, threshold float64, repeatInterval time.Duration, limit int) (*relay.HostAccountLimitAlertClaims, error) { f.lastLimit = limit - return f.usages, nil + if len(f.claims) == 0 { + return &relay.HostAccountLimitAlertClaims{}, nil + } + idx := f.calls + if idx >= len(f.claims) { + idx = len(f.claims) - 1 + } + f.calls++ + return &f.claims[idx], nil +} + +func (f *fakeAccountLimitAlertClaimer) RecordHostAccountLimitAlertSent(ctx context.Context, hostname string, state models.HostAccountLimitAlertState, sentAt time.Time) error { + f.records = append(f.records, hostname+":"+string(state)) + return nil } type recordingAlerter struct { @@ -35,10 +51,14 @@ func (r *recordingAlerter) SendAlert(ctx context.Context, msg AlertMessage) erro return nil } -func TestAccountLimitAlertMonitorSendsAndRepeatsAfterCooldown(t *testing.T) { - lister := &fakeAccountLimitUsageLister{ - usages: []relay.HostAccountLimitUsage{ - {ID: 1, Hostname: "pds.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8}, +func TestAccountLimitAlertMonitorSendsClaimedWarning(t *testing.T) { + claimer := &fakeAccountLimitAlertClaimer{ + claims: []relay.HostAccountLimitAlertClaims{ + { + Warnings: []relay.HostAccountLimitUsage{ + {ID: 1, Hostname: "pds.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8, AlertSentAt: time.Unix(1000, 0)}, + }, + }, }, } alerter := &recordingAlerter{} @@ -47,24 +67,13 @@ func TestAccountLimitAlertMonitorSendsAndRepeatsAfterCooldown(t *testing.T) { config.RepeatInterval = time.Hour config.Environment = "test" - monitor, err := NewAccountLimitAlertMonitor(slog.Default(), lister, alerter, config) + monitor, err := NewAccountLimitAlertMonitor(slog.Default(), claimer, alerter, config) require.NoError(t, err) - now := time.Unix(1000, 0) - monitor.now = func() time.Time { return now } - require.NoError(t, monitor.Check(context.Background())) require.Len(t, alerter.messages, 1) assert.Contains(t, alerter.messages[0].Title, "in test") assert.Contains(t, alerter.messages[0].Text, "`pds.example.com`: *80.0%* used (`80 / 100` repos)") - - now = now.Add(30 * time.Minute) - require.NoError(t, monitor.Check(context.Background())) - require.Len(t, alerter.messages, 1) - - now = now.Add(31 * time.Minute) - require.NoError(t, monitor.Check(context.Background())) - require.Len(t, alerter.messages, 2) } func TestJitteredCheckInterval(t *testing.T) { @@ -106,9 +115,13 @@ func TestInitialCheckDelay(t *testing.T) { } func TestAccountLimitAlertMonitorCirculatesLocalSentState(t *testing.T) { - lister := &fakeAccountLimitUsageLister{ - usages: []relay.HostAccountLimitUsage{ - {ID: 1, Hostname: "pds.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8}, + claimer := &fakeAccountLimitAlertClaimer{ + claims: []relay.HostAccountLimitAlertClaims{ + { + Warnings: []relay.HostAccountLimitUsage{ + {ID: 1, Hostname: "pds.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8, AlertSentAt: time.Unix(1000, 0)}, + }, + }, }, } alerter := &recordingAlerter{} @@ -118,7 +131,7 @@ func TestAccountLimitAlertMonitorCirculatesLocalSentState(t *testing.T) { sent = append(sent, state) } - monitor, err := NewAccountLimitAlertMonitor(slog.Default(), lister, alerter, config) + monitor, err := NewAccountLimitAlertMonitor(slog.Default(), claimer, alerter, config) require.NoError(t, err) require.NoError(t, monitor.Check(context.Background())) @@ -128,111 +141,49 @@ func TestAccountLimitAlertMonitorCirculatesLocalSentState(t *testing.T) { assert.Equal(t, "pds.example.com", sent[0].Hostname) } -func TestAccountLimitAlertMonitorRemoteWarningSuppressesLocalWarning(t *testing.T) { - lister := &fakeAccountLimitUsageLister{ - usages: []relay.HostAccountLimitUsage{ - {ID: 1, Hostname: "pds.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8}, - }, - } - alerter := &recordingAlerter{} - config := DefaultAccountLimitAlertMonitorConfig() - - monitor, err := NewAccountLimitAlertMonitor(slog.Default(), lister, alerter, config) - require.NoError(t, err) - require.NoError(t, monitor.RecordAccountLimitAlertSent(AccountLimitAlertSentState{ - Kind: AccountLimitAlertKindWarning, - HostID: 1, - Hostname: "pds.example.com", - AccountCount: 80, - AccountLimit: 100, - Usage: 0.8, - SentAt: time.Now(), - })) - - require.NoError(t, monitor.Check(context.Background())) - require.Empty(t, alerter.messages) -} - -func TestAccountLimitAlertMonitorRemoteRecoverySuppressesLocalRecovery(t *testing.T) { - lister := &fakeAccountLimitUsageLister{} - alerter := &recordingAlerter{} - config := DefaultAccountLimitAlertMonitorConfig() - - monitor, err := NewAccountLimitAlertMonitor(slog.Default(), lister, alerter, config) - require.NoError(t, err) - require.NoError(t, monitor.RecordAccountLimitAlertSent(AccountLimitAlertSentState{ - Kind: AccountLimitAlertKindWarning, - HostID: 1, - Hostname: "pds.example.com", - AccountCount: 80, - AccountLimit: 100, - Usage: 0.8, - SentAt: time.Now(), - })) - require.NoError(t, monitor.RecordAccountLimitAlertSent(AccountLimitAlertSentState{ - Kind: AccountLimitAlertKindRecovery, - HostID: 1, - Hostname: "pds.example.com", - AccountCount: 80, - AccountLimit: 100, - Usage: 0.8, - SentAt: time.Now(), - })) - - require.NoError(t, monitor.Check(context.Background())) - require.Empty(t, alerter.messages) -} - -func TestAccountLimitAlertMonitorSendsRecoveryAndAlertsAgain(t *testing.T) { - lister := &fakeAccountLimitUsageLister{ - usages: []relay.HostAccountLimitUsage{ - {ID: 1, Hostname: "pds.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8}, +func TestAccountLimitAlertMonitorSendsClaimedRecovery(t *testing.T) { + claimer := &fakeAccountLimitAlertClaimer{ + claims: []relay.HostAccountLimitAlertClaims{ + { + Recoveries: []relay.HostAccountLimitUsage{ + {ID: 1, Hostname: "pds.example.com", AccountCount: 20, AccountLimit: 100, Usage: 0.2, AlertSentAt: time.Unix(1000, 0)}, + }, + }, }, } alerter := &recordingAlerter{} config := DefaultAccountLimitAlertMonitorConfig() config.RepeatInterval = time.Hour - monitor, err := NewAccountLimitAlertMonitor(slog.Default(), lister, alerter, config) + monitor, err := NewAccountLimitAlertMonitor(slog.Default(), claimer, alerter, config) require.NoError(t, err) - now := time.Unix(1000, 0) - monitor.now = func() time.Time { return now } require.NoError(t, monitor.Check(context.Background())) require.Len(t, alerter.messages, 1) - - lister.usages = nil - now = now.Add(10 * time.Minute) - require.NoError(t, monitor.Check(context.Background())) - require.Len(t, alerter.messages, 2) - assert.Contains(t, alerter.messages[1].Title, "PDS repo limit recovered") - assert.Contains(t, alerter.messages[1].Text, "`pds.example.com` is no longer above the 80.0% repo-limit alert threshold") - - lister.usages = []relay.HostAccountLimitUsage{ - {ID: 1, Hostname: "pds.example.com", AccountCount: 81, AccountLimit: 100, Usage: 0.81}, - } - now = now.Add(10 * time.Minute) - require.NoError(t, monitor.Check(context.Background())) - require.Len(t, alerter.messages, 3) - assert.Contains(t, alerter.messages[2].Text, "`pds.example.com`: *81.0%* used (`81 / 100` repos)") + assert.Contains(t, alerter.messages[0].Title, "PDS repo limit recovered") + assert.Contains(t, alerter.messages[0].Text, "`pds.example.com` is no longer above the 80.0% repo-limit alert threshold") } func TestAccountLimitAlertMonitorBatchesMessagesWithoutQueryCap(t *testing.T) { - lister := &fakeAccountLimitUsageLister{ - usages: []relay.HostAccountLimitUsage{ - {ID: 1, Hostname: "one.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8}, - {ID: 2, Hostname: "two.example.com", AccountCount: 81, AccountLimit: 100, Usage: 0.81}, + claimer := &fakeAccountLimitAlertClaimer{ + claims: []relay.HostAccountLimitAlertClaims{ + { + Warnings: []relay.HostAccountLimitUsage{ + {ID: 1, Hostname: "one.example.com", AccountCount: 80, AccountLimit: 100, Usage: 0.8, AlertSentAt: time.Unix(1000, 0)}, + {ID: 2, Hostname: "two.example.com", AccountCount: 81, AccountLimit: 100, Usage: 0.81, AlertSentAt: time.Unix(1000, 0)}, + }, + }, }, } alerter := &recordingAlerter{} config := DefaultAccountLimitAlertMonitorConfig() config.HostsPerMessage = 1 - monitor, err := NewAccountLimitAlertMonitor(slog.Default(), lister, alerter, config) + monitor, err := NewAccountLimitAlertMonitor(slog.Default(), claimer, alerter, config) require.NoError(t, err) require.NoError(t, monitor.Check(context.Background())) - assert.Equal(t, 0, lister.lastLimit) + assert.Equal(t, 0, claimer.lastLimit) require.Len(t, alerter.messages, 2) assert.Contains(t, alerter.messages[0].Text, "one.example.com") assert.Contains(t, alerter.messages[1].Text, "two.example.com") diff --git a/cmd/relay/handlers_admin.go b/cmd/relay/handlers_admin.go index 3190ee69..35d6b1af 100644 --- a/cmd/relay/handlers_admin.go +++ b/cmd/relay/handlers_admin.go @@ -219,6 +219,10 @@ type hostInfo struct { PerHourEventRate rateLimit `json:"PerHourEventRate"` PerDayEventRate rateLimit `json:"PerDayEventRate"` UserCount int64 `json:"UserCount"` + + AccountLimitAlertsSilenced bool `json:"AccountLimitAlertsSilenced"` + AccountLimitAlertState string `json:"AccountLimitAlertState"` + AccountLimitAlertSentAt *time.Time `json:"AccountLimitAlertSentAt,omitempty"` } func (s *Service) handleListHosts(c echo.Context) error { @@ -252,6 +256,10 @@ func (s *Service) handleListHosts(c echo.Context) error { HasActiveConnection: isActive, UserCount: host.AccountCount, + + AccountLimitAlertsSilenced: host.AccountLimitAlertsSilenced, + AccountLimitAlertState: string(host.AccountLimitAlertState), + AccountLimitAlertSentAt: host.AccountLimitAlertSentAt, } // fetch current rate limits @@ -470,8 +478,9 @@ func (s *Service) handleAdminUnbanDomain(c echo.Context) error { } type RateLimitChangeRequest struct { - Hostname string `json:"host"` - RepoLimit *int64 `json:"repo_limit"` + Hostname string `json:"host"` + RepoLimit *int64 `json:"repo_limit"` + AccountLimitAlertsSilenced *bool `json:"account_limit_alerts_silenced"` } func (s *Service) handleAdminChangeHostRateLimits(c echo.Context) error { @@ -487,9 +496,8 @@ func (s *Service) handleAdminChangeHostRateLimits(c echo.Context) error { return echo.NewHTTPError(http.StatusBadRequest, fmt.Sprintf("invalid hostname: %s", err)) } - // catch empty/nil body - if body.RepoLimit == nil { - return echo.NewHTTPError(http.StatusBadRequest, "missing repo_limit parameter") + if body.RepoLimit == nil && body.AccountLimitAlertsSilenced == nil { + return echo.NewHTTPError(http.StatusBadRequest, "missing repo_limit or account_limit_alerts_silenced parameter") } host, err := s.relay.GetHost(ctx, hostname) @@ -498,8 +506,15 @@ func (s *Service) handleAdminChangeHostRateLimits(c echo.Context) error { return echo.NewHTTPError(http.StatusNotFound, fmt.Sprintf("unknown hostname: %s", err)) } - if err := s.relay.UpdateHostAccountLimit(ctx, host.ID, *body.RepoLimit); err != nil { - return echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("failed to update limits: %s", err)) + if body.RepoLimit != nil { + if err := s.relay.UpdateHostAccountLimit(ctx, host.ID, *body.RepoLimit); err != nil { + return echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("failed to update limits: %s", err)) + } + } + if body.AccountLimitAlertsSilenced != nil { + if err := s.relay.UpdateHostAccountLimitAlertsSilenced(ctx, host.ID, *body.AccountLimitAlertsSilenced); err != nil { + return echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("failed to update account limit alert silence state: %s", err)) + } } // forward on to any sibling instances diff --git a/cmd/relay/main.go b/cmd/relay/main.go index f6ed7395..7376293e 100644 --- a/cmd/relay/main.go +++ b/cmd/relay/main.go @@ -391,9 +391,6 @@ func runRelay(ctx context.Context, cmd *cli.Command) error { if err != nil { return err } - if alertMonitor != nil { - svc.SetAccountLimitAlertRecorder(alertMonitor) - } alertCtx, cancelAlerts := context.WithCancel(ctx) alertDone := make(chan struct{}) if alertMonitor != nil { diff --git a/cmd/relay/relay-admin-ui/src/components/Dash/Dash.tsx b/cmd/relay/relay-admin-ui/src/components/Dash/Dash.tsx index dc787b6c..61b795d7 100644 --- a/cmd/relay/relay-admin-ui/src/components/Dash/Dash.tsx +++ b/cmd/relay/relay-admin-ui/src/components/Dash/Dash.tsx @@ -412,6 +412,55 @@ const Dash: FC<{}> = () => { }); }; + const updateAccountLimitAlertsSilenced = (pds: PDS, silenced: boolean) => { + const applyOptimisticSilenceState = (list: PDS[] | null): PDS[] | null => { + if (!list) { + return list; + } + return list.map((item) => { + if (item.ID !== pds.ID) { + return item; + } + return { + ...item, + AccountLimitAlertsSilenced: silenced, + AccountLimitAlertState: silenced ? "ok" : item.AccountLimitAlertState, + AccountLimitAlertSentAt: silenced ? undefined : item.AccountLimitAlertSentAt, + }; + }); + }; + + setPDSList(applyOptimisticSilenceState); + setFullPDSList(applyOptimisticSilenceState); + + fetch(`${RELAY_HOST}/admin/pds/changeLimits`, { + method: "POST", + headers: { + "Content-Type": "application/json", + Authorization: `Basic ` + btoa("admin:" + adminToken), + }, + body: JSON.stringify({ + host: pds.Host, + account_limit_alerts_silenced: silenced, + }), + }).then((res) => { + if (res.status !== 200) { + setAlertWithTimeout( + "failure", + `Failed to ${silenced ? "silence" : "unsilence"} account limit alerts: ${res.statusText} (${res.status})`, + true + ); + } else { + setAlertWithTimeout( + "success", + `Successfully ${silenced ? "silenced" : "unsilenced"} account limit alerts`, + true + ); + } + refreshPDSList(); + }); + }; + const handleBlockClick = (pds: PDS, shouldBlock: boolean) => { setModalAction({ type: shouldBlock ? "block" : "disconnect", @@ -860,6 +909,14 @@ const Dash: FC<{}> = () => { Repo Limit + + + Alerts + + = () => { /> + + + {new Date(Date.parse(pds.CreatedAt)).toLocaleString()} diff --git a/cmd/relay/relay-admin-ui/src/models/pds.ts b/cmd/relay/relay-admin-ui/src/models/pds.ts index 64b1135d..3421430e 100644 --- a/cmd/relay/relay-admin-ui/src/models/pds.ts +++ b/cmd/relay/relay-admin-ui/src/models/pds.ts @@ -21,6 +21,9 @@ interface PDS { PerDayEventRate: RateLimit; RepoCount: number; RepoLimit: number; + AccountLimitAlertsSilenced: boolean; + AccountLimitAlertState: string; + AccountLimitAlertSentAt?: string; } type PDSKey = keyof PDS; diff --git a/cmd/relay/relay/host_usage.go b/cmd/relay/relay/host_usage.go index 7f00bcd9..0d7ac19a 100644 --- a/cmd/relay/relay/host_usage.go +++ b/cmd/relay/relay/host_usage.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "sort" + "time" "github.com/bluesky-social/indigo/cmd/relay/relay/models" ) @@ -14,6 +15,12 @@ type HostAccountLimitUsage struct { AccountCount int64 AccountLimit int64 Usage float64 + AlertSentAt time.Time +} + +type HostAccountLimitAlertClaims struct { + Warnings []HostAccountLimitUsage + Recoveries []HostAccountLimitUsage } func (r *Relay) ListHostsApproachingAccountLimit(ctx context.Context, threshold float64, limit int) ([]HostAccountLimitUsage, error) { @@ -56,3 +63,148 @@ func (r *Relay) ListHostsApproachingAccountLimit(ctx context.Context, threshold return out, nil } + +func (r *Relay) ClaimDueAccountLimitAlerts(ctx context.Context, threshold float64, repeatInterval time.Duration, limit int) (*HostAccountLimitAlertClaims, error) { + if threshold <= 0 || threshold > 1 { + return nil, fmt.Errorf("account limit alert threshold must be greater than 0 and less than or equal to 1") + } + if repeatInterval <= 0 { + return nil, fmt.Errorf("account limit alert repeat interval must be positive") + } + if limit < 0 { + return nil, fmt.Errorf("account limit alert query limit must not be negative") + } + + now := time.Now().UTC() + cutoff := now.Add(-repeatInterval) + + warnings, err := r.claimDueAccountLimitWarnings(ctx, threshold, cutoff, now, limit) + if err != nil { + return nil, err + } + recoveries, err := r.claimDueAccountLimitRecoveries(ctx, threshold, now, limit) + if err != nil { + return nil, err + } + return &HostAccountLimitAlertClaims{ + Warnings: warnings, + Recoveries: recoveries, + }, nil +} + +func (r *Relay) claimDueAccountLimitWarnings(ctx context.Context, threshold float64, cutoff time.Time, now time.Time, limit int) ([]HostAccountLimitUsage, error) { + var hosts []models.Host + query := r.db.WithContext(ctx). + Model(&models.Host{}). + Where("account_limit > 0 AND account_count >= account_limit * ? AND status <> ? AND account_limit_alerts_silenced = ?", threshold, models.HostStatusBanned, false). + Where("(account_limit_alert_state <> ? OR account_limit_alert_state IS NULL OR account_limit_alert_sent_at IS NULL OR account_limit_alert_sent_at <= ?)", models.HostAccountLimitAlertStateWarning, cutoff). + Order("account_count * 1.0 / account_limit DESC, id ASC") + if limit > 0 { + query = query.Limit(limit) + } + if err := query.Find(&hosts).Error; err != nil { + return nil, err + } + + out := make([]HostAccountLimitUsage, 0, len(hosts)) + for _, host := range hosts { + result := r.db.WithContext(ctx).Model(&models.Host{}). + Where("id = ? AND account_limit > 0 AND account_count >= account_limit * ? AND status <> ? AND account_limit_alerts_silenced = ?", host.ID, threshold, models.HostStatusBanned, false). + Where("(account_limit_alert_state <> ? OR account_limit_alert_state IS NULL OR account_limit_alert_sent_at IS NULL OR account_limit_alert_sent_at <= ?)", models.HostAccountLimitAlertStateWarning, cutoff). + Updates(map[string]any{ + "account_limit_alert_state": models.HostAccountLimitAlertStateWarning, + "account_limit_alert_sent_at": now, + }) + if result.Error != nil { + return nil, result.Error + } + if result.RowsAffected == 0 { + continue + } + out = append(out, hostAccountLimitUsage(host, now)) + } + return out, nil +} + +func (r *Relay) claimDueAccountLimitRecoveries(ctx context.Context, threshold float64, now time.Time, limit int) ([]HostAccountLimitUsage, error) { + var hosts []models.Host + query := r.db.WithContext(ctx). + Model(&models.Host{}). + Where("account_limit > 0 AND account_count < account_limit * ? AND status <> ? AND account_limit_alerts_silenced = ? AND account_limit_alert_state = ?", threshold, models.HostStatusBanned, false, models.HostAccountLimitAlertStateWarning). + Order("account_count * 1.0 / account_limit DESC, id ASC") + if limit > 0 { + query = query.Limit(limit) + } + if err := query.Find(&hosts).Error; err != nil { + return nil, err + } + + out := make([]HostAccountLimitUsage, 0, len(hosts)) + for _, host := range hosts { + result := r.db.WithContext(ctx).Model(&models.Host{}). + Where("id = ? AND account_limit > 0 AND account_count < account_limit * ? AND status <> ? AND account_limit_alerts_silenced = ? AND account_limit_alert_state = ?", host.ID, threshold, models.HostStatusBanned, false, models.HostAccountLimitAlertStateWarning). + Updates(map[string]any{ + "account_limit_alert_state": models.HostAccountLimitAlertStateOK, + "account_limit_alert_sent_at": now, + }) + if result.Error != nil { + return nil, result.Error + } + if result.RowsAffected == 0 { + continue + } + out = append(out, hostAccountLimitUsage(host, now)) + } + return out, nil +} + +func (r *Relay) RecordHostAccountLimitAlertSent(ctx context.Context, hostname string, state models.HostAccountLimitAlertState, sentAt time.Time) error { + switch state { + case models.HostAccountLimitAlertStateOK, models.HostAccountLimitAlertStateWarning: + default: + return fmt.Errorf("invalid account limit alert state: %s", state) + } + if sentAt.IsZero() { + sentAt = time.Now().UTC() + } + result := r.db.WithContext(ctx).Model(&models.Host{}). + Where("hostname = ? AND account_limit_alerts_silenced = ?", hostname, false). + Where("account_limit_alert_sent_at IS NULL OR account_limit_alert_sent_at <= ?", sentAt). + Updates(map[string]any{ + "account_limit_alert_state": state, + "account_limit_alert_sent_at": sentAt, + }) + if result.Error != nil { + return result.Error + } + return nil +} + +func (r *Relay) UpdateHostAccountLimitAlertsSilenced(ctx context.Context, hostID uint64, silenced bool) error { + updates := map[string]any{ + "account_limit_alerts_silenced": silenced, + } + if silenced { + updates["account_limit_alert_state"] = models.HostAccountLimitAlertStateOK + updates["account_limit_alert_sent_at"] = nil + } + result := r.db.WithContext(ctx).Model(models.Host{}).Where("id = ?", hostID).Updates(updates) + if result.Error != nil { + return result.Error + } + if result.RowsAffected == 0 { + return ErrHostNotFound + } + return nil +} + +func hostAccountLimitUsage(h models.Host, alertSentAt time.Time) HostAccountLimitUsage { + return HostAccountLimitUsage{ + ID: h.ID, + Hostname: h.Hostname, + AccountCount: h.AccountCount, + AccountLimit: h.AccountLimit, + Usage: float64(h.AccountCount) / float64(h.AccountLimit), + AlertSentAt: alertSentAt, + } +} diff --git a/cmd/relay/relay/host_usage_test.go b/cmd/relay/relay/host_usage_test.go index 9b346910..1659f7df 100644 --- a/cmd/relay/relay/host_usage_test.go +++ b/cmd/relay/relay/host_usage_test.go @@ -3,6 +3,7 @@ package relay import ( "context" "testing" + "time" "github.com/bluesky-social/indigo/cmd/relay/relay/models" @@ -12,10 +13,16 @@ import ( "gorm.io/gorm" ) -func TestListHostsApproachingAccountLimit(t *testing.T) { +func testRelayWithHostDB(t *testing.T) (*Relay, *gorm.DB) { + t.Helper() db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) require.NoError(t, err) require.NoError(t, db.AutoMigrate(&models.Host{})) + return &Relay{db: db}, db +} + +func TestListHostsApproachingAccountLimit(t *testing.T) { + r, db := testRelayWithHostDB(t) hosts := []models.Host{ {Hostname: "below.example.com", AccountCount: 79, AccountLimit: 100, Status: models.HostStatusActive}, @@ -27,7 +34,6 @@ func TestListHostsApproachingAccountLimit(t *testing.T) { require.NoError(t, db.Create(&host).Error) } - r := &Relay{db: db} usages, err := r.ListHostsApproachingAccountLimit(context.Background(), 0.80, 10) require.NoError(t, err) require.Len(t, usages, 2) @@ -39,6 +45,139 @@ func TestListHostsApproachingAccountLimit(t *testing.T) { assert.Equal(t, "at.example.com", usages[1].Hostname) } +func TestClaimDueAccountLimitAlertsPersistsWarningState(t *testing.T) { + ctx := context.Background() + r, db := testRelayWithHostDB(t) + require.NoError(t, db.Create(&models.Host{ + Hostname: "over.example.com", + AccountCount: 85, + AccountLimit: 100, + Status: models.HostStatusActive, + }).Error) + + claims, err := r.ClaimDueAccountLimitAlerts(ctx, 0.8, time.Hour, 0) + require.NoError(t, err) + require.Len(t, claims.Warnings, 1) + require.Empty(t, claims.Recoveries) + assert.Equal(t, "over.example.com", claims.Warnings[0].Hostname) + assert.False(t, claims.Warnings[0].AlertSentAt.IsZero()) + + claims, err = r.ClaimDueAccountLimitAlerts(ctx, 0.8, time.Hour, 0) + require.NoError(t, err) + require.Empty(t, claims.Warnings) + require.Empty(t, claims.Recoveries) + + var host models.Host + require.NoError(t, db.First(&host, "hostname = ?", "over.example.com").Error) + assert.Equal(t, models.HostAccountLimitAlertStateWarning, host.AccountLimitAlertState) + require.NotNil(t, host.AccountLimitAlertSentAt) +} + +func TestClaimDueAccountLimitAlertsRepeatsAfterInterval(t *testing.T) { + ctx := context.Background() + r, db := testRelayWithHostDB(t) + oldSentAt := time.Now().Add(-2 * time.Hour) + require.NoError(t, db.Create(&models.Host{ + Hostname: "over.example.com", + AccountCount: 85, + AccountLimit: 100, + Status: models.HostStatusActive, + AccountLimitAlertState: models.HostAccountLimitAlertStateWarning, + AccountLimitAlertSentAt: &oldSentAt, + }).Error) + + claims, err := r.ClaimDueAccountLimitAlerts(ctx, 0.8, time.Hour, 0) + require.NoError(t, err) + require.Len(t, claims.Warnings, 1) + require.Empty(t, claims.Recoveries) +} + +func TestClaimDueAccountLimitAlertsClaimsRecovery(t *testing.T) { + ctx := context.Background() + r, db := testRelayWithHostDB(t) + sentAt := time.Now().Add(-time.Hour) + require.NoError(t, db.Create(&models.Host{ + Hostname: "recovered.example.com", + AccountCount: 20, + AccountLimit: 100, + Status: models.HostStatusActive, + AccountLimitAlertState: models.HostAccountLimitAlertStateWarning, + AccountLimitAlertSentAt: &sentAt, + }).Error) + + claims, err := r.ClaimDueAccountLimitAlerts(ctx, 0.8, time.Hour, 0) + require.NoError(t, err) + require.Empty(t, claims.Warnings) + require.Len(t, claims.Recoveries, 1) + assert.Equal(t, "recovered.example.com", claims.Recoveries[0].Hostname) + + var host models.Host + require.NoError(t, db.First(&host, "hostname = ?", "recovered.example.com").Error) + assert.Equal(t, models.HostAccountLimitAlertStateOK, host.AccountLimitAlertState) +} + +func TestClaimDueAccountLimitAlertsSkipsSilencedHosts(t *testing.T) { + ctx := context.Background() + r, db := testRelayWithHostDB(t) + require.NoError(t, db.Create(&models.Host{ + Hostname: "silenced.example.com", + AccountCount: 95, + AccountLimit: 100, + Status: models.HostStatusActive, + AccountLimitAlertsSilenced: true, + }).Error) + + claims, err := r.ClaimDueAccountLimitAlerts(ctx, 0.8, time.Hour, 0) + require.NoError(t, err) + require.Empty(t, claims.Warnings) + require.Empty(t, claims.Recoveries) +} + +func TestUpdateHostAccountLimitAlertsSilencedResetsAlertState(t *testing.T) { + ctx := context.Background() + r, db := testRelayWithHostDB(t) + sentAt := time.Now().Add(-time.Hour) + host := models.Host{ + Hostname: "silence-me.example.com", + AccountCount: 95, + AccountLimit: 100, + Status: models.HostStatusActive, + AccountLimitAlertState: models.HostAccountLimitAlertStateWarning, + AccountLimitAlertSentAt: &sentAt, + } + require.NoError(t, db.Create(&host).Error) + + require.NoError(t, r.UpdateHostAccountLimitAlertsSilenced(ctx, host.ID, true)) + + var updated models.Host + require.NoError(t, db.First(&updated, host.ID).Error) + assert.True(t, updated.AccountLimitAlertsSilenced) + assert.Equal(t, models.HostAccountLimitAlertStateOK, updated.AccountLimitAlertState) + assert.Nil(t, updated.AccountLimitAlertSentAt) +} + +func TestRecordHostAccountLimitAlertSentIgnoresStaleState(t *testing.T) { + ctx := context.Background() + r, db := testRelayWithHostDB(t) + newerSentAt := time.Now() + require.NoError(t, db.Create(&models.Host{ + Hostname: "state.example.com", + AccountCount: 95, + AccountLimit: 100, + Status: models.HostStatusActive, + AccountLimitAlertState: models.HostAccountLimitAlertStateWarning, + AccountLimitAlertSentAt: &newerSentAt, + }).Error) + + require.NoError(t, r.RecordHostAccountLimitAlertSent(ctx, "state.example.com", models.HostAccountLimitAlertStateOK, newerSentAt.Add(-time.Minute))) + + var host models.Host + require.NoError(t, db.First(&host, "hostname = ?", "state.example.com").Error) + assert.Equal(t, models.HostAccountLimitAlertStateWarning, host.AccountLimitAlertState) + require.NotNil(t, host.AccountLimitAlertSentAt) + assert.True(t, host.AccountLimitAlertSentAt.Equal(newerSentAt)) +} + func TestListHostsApproachingAccountLimitRejectsInvalidThreshold(t *testing.T) { r := &Relay{} _, err := r.ListHostsApproachingAccountLimit(context.Background(), 0, 10) diff --git a/cmd/relay/relay/models/models.go b/cmd/relay/relay/models/models.go index ae76f336..43f8b431 100644 --- a/cmd/relay/relay/models/models.go +++ b/cmd/relay/relay/models/models.go @@ -22,6 +22,13 @@ const ( HostStatusBanned = HostStatus("banned") ) +type HostAccountLimitAlertState string + +const ( + HostAccountLimitAlertStateOK = HostAccountLimitAlertState("ok") + HostAccountLimitAlertStateWarning = HostAccountLimitAlertState("warning") +) + type Host struct { ID uint64 `gorm:"column:id;primarykey" json:"id"` @@ -48,6 +55,15 @@ type Host struct { // represents the number of accounts on the host, minus any in "deleted" state AccountCount int64 `gorm:"column:account_count;default:0" json:"accountCount"` + + // if true, account-limit alerting is disabled for this host + AccountLimitAlertsSilenced bool `gorm:"column:account_limit_alerts_silenced;default:false" json:"accountLimitAlertsSilenced"` + + // persisted account-limit alert state, used to avoid duplicate alerts across process restarts + AccountLimitAlertState HostAccountLimitAlertState `gorm:"column:account_limit_alert_state;default:ok" json:"accountLimitAlertState"` + + // last account-limit warning or recovery alert send time + AccountLimitAlertSentAt *time.Time `gorm:"column:account_limit_alert_sent_at" json:"accountLimitAlertSentAt,omitempty"` } func (Host) TableName() string { diff --git a/cmd/relay/service.go b/cmd/relay/service.go index 15194b86..cfdf8235 100644 --- a/cmd/relay/service.go +++ b/cmd/relay/service.go @@ -24,12 +24,6 @@ type Service struct { config ServiceConfig siblingClient http.Client - - accountLimitAlertRecorder accountLimitAlertRecorder -} - -type accountLimitAlertRecorder interface { - RecordAccountLimitAlertSent(state AccountLimitAlertSentState) error } type ServiceConfig struct { @@ -73,10 +67,6 @@ func NewService(r *relay.Relay, config *ServiceConfig) (*Service, error) { return svc, nil } -func (svc *Service) SetAccountLimitAlertRecorder(recorder accountLimitAlertRecorder) { - svc.accountLimitAlertRecorder = recorder -} - func (svc *Service) StartMetrics(listen string) error { http.Handle("/metrics", promhttp.Handler()) return http.ListenAndServe(listen, nil)