feat(instance): метрики и healthcheck инстанса для дашборда

Раздел «Инстанс» получает живые показатели сервера (AGENT.md 7.19).

- `internal/sysinfo`: загрузка процессора по cgroup v2 (`cpu.stat`/`cpu.max`)
  с откатом на `/proc/stat`, память по лимиту контейнера либо `/proc/meminfo`,
  место на разделе данных и размер базы вместе с журналом WAL;
- `GET /api/v1/instance/metrics` (только администратор инстанса): метрики,
  состояние (`ok|warn|failed`), проверки app/database/disk/storage/voice/
  gateway, версия и время работы; недоступные источники не ломают ответ, а
  попадают в проверку `metrics`;
- тяжёлые проверки (БД, SFU, запись) кэшируются на 5 секунд: ручку можно
  опрашивать раз в секунду;
- `voice.IssueRoomList` и `AdminClient.Ping` проверяют доступность RoomService
  и валидность ключей LiveKit (та же связка, что ломала модерацию);
- у метрик свой лимит частоты, они исключены из общего лимита API
  (`RateLimiter.MiddlewareExcept`), иначе открытый раздел съедал бы половину
  бюджета запросов;
- тесты: `internal/sysinfo` (cgroup, /proc, память, диск, размеры файлов) и
  `internal/server/api_metrics_test.go` (доступ, состав метрик, кэш, 429).
This commit is contained in:
2026-09-20 21:26:17 +03:00
parent 8ae8ca1a1d
commit 4c3736b5fb
8 changed files with 1128 additions and 19 deletions
+21
View File
@@ -124,6 +124,27 @@ func (l *RateLimiter) Middleware(keyFn func(*http.Request) string) Middleware {
} }
} }
// MiddlewareExcept ограничивает частоту запросов везде, кроме перечисленных
// путей: у живого дашборда инстанса свой лимит, иначе опрос раз в секунду
// съедал бы половину общего бюджета API (AGENT.md 8.6).
func (l *RateLimiter) MiddlewareExcept(keyFn func(*http.Request) string, except ...string) Middleware {
skip := make(map[string]struct{}, len(except))
for _, path := range except {
skip[path] = struct{}{}
}
limited := l.Middleware(keyFn)
return func(next http.Handler) http.Handler {
handler := limited(next)
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if _, ok := skip[r.URL.Path]; ok {
next.ServeHTTP(w, r)
return
}
handler.ServeHTTP(w, r)
})
}
}
// WriteRateLimited отдаёт 429 в конверте API с подсказкой по паузе (§8.5, §8.6). // WriteRateLimited отдаёт 429 в конверте API с подсказкой по паузе (§8.5, §8.6).
func WriteRateLimited(w http.ResponseWriter, retryAfter time.Duration) { func WriteRateLimited(w http.ResponseWriter, retryAfter time.Duration) {
milliseconds := retryAfter.Milliseconds() milliseconds := retryAfter.Milliseconds()
+313
View File
@@ -0,0 +1,313 @@
package server
import (
"context"
"fmt"
"net/http"
"os"
"path/filepath"
"strings"
"sync"
"time"
"github.com/danielgtaylor/huma/v2"
"glchat/internal/sysinfo"
)
// metricsPath — ручка живого дашборда инстанса. Вынесена в константу: её
// опрос раз в секунду исключён из общего лимита API (см. registerRoutes).
const metricsPath = "/api/v1/instance/metrics"
// healthCacheTTL — как долго переиспользуются результаты тяжёлых проверок
// (запрос к БД и к SFU). Метрики при этом считаются на каждый запрос, поэтому
// ручку можно опрашивать раз в секунду.
const healthCacheTTL = 5 * time.Second
// healthCache хранит последний healthcheck: дашборд обновляется раз в секунду,
// а проверки с сетевыми вызовами чаще раза в healthCacheTTL не нужны.
type healthCache struct {
mu sync.Mutex
checks []metricsCheck
at time.Time
}
type metricsCPU struct {
Percent float64 `json:"percent" doc:"Загрузка процессора в процентах от доступного времени"`
Cores float64 `json:"cores" doc:"Доступные ядра (лимит контейнера или ядра хоста)"`
Source string `json:"source" doc:"Источник данных: cgroup или host"`
}
type metricsMemory struct {
UsedBytes int64 `json:"used_bytes"`
TotalBytes int64 `json:"total_bytes"`
Percent float64 `json:"percent"`
Source string `json:"source" doc:"Источник данных: cgroup (лимит) или host"`
}
type metricsDisk struct {
UsedBytes int64 `json:"used_bytes"`
TotalBytes int64 `json:"total_bytes"`
FreeBytes int64 `json:"free_bytes"`
}
type metricsDatabase struct {
Bytes int64 `json:"bytes" doc:"Размер файла базы"`
WalBytes int64 `json:"wal_bytes" doc:"Размер журнала WAL"`
}
// metricsCheck — одна проверка инстанса. name и status машиночитаемые, текст
// сообщения технический (локализуется на клиенте по name).
type metricsCheck struct {
Name string `json:"name"`
Status string `json:"status" enum:"ok,warn,failed,off"`
Message string `json:"message,omitempty"`
}
type metricsPayload struct {
CPU metricsCPU `json:"cpu"`
Memory metricsMemory `json:"memory"`
Disk metricsDisk `json:"disk"`
Database metricsDatabase `json:"database"`
Checks []metricsCheck `json:"checks"`
Health string `json:"health" enum:"ok,warn,failed"`
UptimeSeconds int64 `json:"uptime_seconds"`
Version string `json:"version"`
Commit string `json:"commit"`
CollectedAt string `json:"collected_at"`
}
type metricsOutput struct {
Body struct {
Metrics metricsPayload `json:"metrics"`
}
}
// registerMetricsRoutes описывает дашборд инстанса: живые метрики, состояние
// диска и базы, healthcheck и версия (AGENT.md 7.19).
func (s *Server) registerMetricsRoutes(api huma.API) {
huma.Register(api, huma.Operation{
OperationID: "getInstanceMetrics",
Method: http.MethodGet,
Path: "/instance/metrics",
Summary: "Живые метрики инстанса (только администратор)",
Tags: []string{"Instance"},
Security: []map[string][]string{{"sessionCookie": {}}, {"bearerAuth": {}}},
}, func(ctx context.Context, _ *struct{}) (*metricsOutput, error) {
user, _, err := requireUser(ctx)
if err != nil {
return nil, err
}
if !user.IsInstanceAdmin {
return nil, humaErrorStatus(http.StatusForbidden, "perm.denied", "instance admin is required")
}
// Свой лимит: раз в секунду — это 60 запросов в минуту, под общим
// лимитом API такой опрос съедал бы половину бюджета (AGENT.md 8.6).
if allowed, retryAfter := s.metricsLimiter.Allow("metrics:" + formatSnowflake(user.ID)); !allowed {
return nil, rateLimitedError(retryAfter)
}
output := &metricsOutput{}
output.Body.Metrics = s.collectMetrics(ctx)
return output, nil
})
}
// collectMetrics собирает метрики и (при необходимости) обновляет healthcheck.
// Недоступные источники не ломают ответ: они попадают в проверку «metrics»,
// а соответствующий блок остаётся пустым (так ведёт себя, например, macOS без
// cgroup и /proc).
func (s *Server) collectMetrics(ctx context.Context) metricsPayload {
payload := metricsPayload{
UptimeSeconds: int64(time.Since(startedAt).Seconds()),
Version: s.cfg.Version,
Commit: s.cfg.Commit,
CollectedAt: time.Now().UTC().Format(time.RFC3339Nano),
}
var problems []string
if s.cpuSampler == nil {
problems = append(problems, "cpu: sampler is not configured")
} else if percent, cores, source, err := s.cpuSampler.CPUPercent(); err != nil {
problems = append(problems, "cpu: "+err.Error())
} else {
payload.CPU = metricsCPU{Percent: round1(percent), Cores: cores, Source: source}
}
if memory, err := sysinfo.ReadMemory(s.sysinfoPaths); err != nil {
problems = append(problems, "memory: "+err.Error())
} else {
payload.Memory = metricsMemory{
UsedBytes: memory.UsedBytes,
TotalBytes: memory.TotalBytes,
Percent: round1(memory.Percent()),
Source: memory.Source,
}
}
if disk, err := sysinfo.ReadDisk(s.cfg.DataDir); err != nil {
problems = append(problems, "disk: "+err.Error())
} else {
payload.Disk = metricsDisk{
UsedBytes: disk.UsedBytes,
TotalBytes: disk.TotalBytes,
FreeBytes: disk.FreeBytes,
}
}
dbPath := s.databasePath()
payload.Database = metricsDatabase{
Bytes: sysinfo.FileBytes(dbPath),
// WAL живёт рядом с базой и в пике бывает больше её самой.
WalBytes: sysinfo.FileBytes(dbPath + "-wal"),
}
payload.Checks = s.instanceChecks(ctx, problems)
payload.Health = summarizeHealth(payload.Checks)
return payload
}
// databasePath возвращает путь к файлу базы: в конфиге он может быть не задан
// (сборка конфига вручную, тесты), тогда берём его у открытой базы.
func (s *Server) databasePath() string {
if s.cfg.DatabasePath != "" {
return s.cfg.DatabasePath
}
if s.db != nil {
return s.db.Path()
}
return ""
}
// instanceChecks отдаёт проверки инстанса из кэша либо пересчитывает их.
func (s *Server) instanceChecks(ctx context.Context, problems []string) []metricsCheck {
s.health.mu.Lock()
defer s.health.mu.Unlock()
if s.health.checks == nil || time.Since(s.health.at) >= healthCacheTTL {
s.health.checks = s.buildChecks(ctx)
s.health.at = time.Now()
}
if len(problems) == 0 {
return s.health.checks
}
// Проблемы с источниками метрик показываем отдельной проверкой, не затирая
// кэш: причина обычно постоянная (нет cgroup), а не разовая.
checks := make([]metricsCheck, 0, len(s.health.checks)+1)
checks = append(checks, s.health.checks...)
checks = append(checks, metricsCheck{
Name: "metrics",
Status: "warn",
Message: strings.Join(problems, "; "),
})
return checks
}
// buildChecks выполняет проверки: приложение, база, диск, SFU, Gateway.
func (s *Server) buildChecks(ctx context.Context) []metricsCheck {
checks := []metricsCheck{{
Name: "app",
Status: "ok",
Message: s.cfg.Version,
}}
// База: ping и версия схемы (миграции применяются при старте).
if s.db != nil {
checkCtx, cancel := context.WithTimeout(ctx, 3*time.Second)
if err := s.db.Ready(checkCtx); err != nil {
checks = append(checks, metricsCheck{Name: "database", Status: "failed", Message: err.Error()})
} else {
version, err := s.db.SchemaVersion(checkCtx)
switch {
case err != nil:
checks = append(checks, metricsCheck{Name: "database", Status: "warn", Message: err.Error()})
default:
checks = append(checks, metricsCheck{
Name: "database",
Status: "ok",
Message: fmt.Sprintf("schema_version=%d", version),
})
}
}
cancel()
} else {
checks = append(checks, metricsCheck{Name: "database", Status: "off"})
}
// Диск: предупреждаем заранее, до заполнения раздела.
if disk, err := sysinfo.ReadDisk(s.cfg.DataDir); err != nil {
checks = append(checks, metricsCheck{Name: "disk", Status: "warn", Message: err.Error()})
} else {
status := "ok"
switch freePercent := float64(disk.FreeBytes) / float64(disk.TotalBytes) * 100; {
case freePercent < 5:
status = "failed"
case freePercent < 15:
status = "warn"
}
checks = append(checks, metricsCheck{
Name: "disk",
Status: status,
Message: fmt.Sprintf("%.1f%% free", float64(disk.FreeBytes)/float64(disk.TotalBytes)*100),
})
}
// Запись в каталог данных: база, вложения и звуки должны быть записываемы.
checks = append(checks, s.checkDataWritable())
// SFU: заодно проверяются ключи LiveKit — их неверность ломает модерацию.
switch {
case !s.voice.Enabled():
checks = append(checks, metricsCheck{Name: "voice", Status: "off"})
case !s.voiceAdmin.Enabled():
checks = append(checks, metricsCheck{Name: "voice", Status: "off", Message: "api url is not set"})
default:
checkCtx, cancel := context.WithTimeout(ctx, 3*time.Second)
defer cancel()
if err := s.voiceAdmin.Ping(checkCtx); err != nil {
checks = append(checks, metricsCheck{Name: "voice", Status: "failed", Message: err.Error()})
} else {
checks = append(checks, metricsCheck{Name: "voice", Status: "ok", Message: "RoomService"})
}
}
if s.gateway != nil {
checks = append(checks, metricsCheck{
Name: "gateway",
Status: "ok",
Message: fmt.Sprintf("sessions=%d", s.gateway.ActiveSessions()),
})
}
return checks
}
// checkDataWritable проверяет, что в каталоге данных можно создавать файлы.
func (s *Server) checkDataWritable() metricsCheck {
probe := filepath.Join(s.cfg.DataDir, ".healthcheck")
// Путь собирается из каталога данных инстанса, а не из пользовательского ввода.
file, err := os.OpenFile(probe, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o600) //nolint:gosec // каталог данных задаёт конфигурация
if err != nil {
return metricsCheck{Name: "storage", Status: "failed", Message: err.Error()}
}
_ = file.Close()
_ = os.Remove(probe)
return metricsCheck{Name: "storage", Status: "ok"}
}
// summarizeHealth сводит проверки к общему состоянию: failed → failed,
// warn/off пропускаем как «не критично», иначе ok.
func summarizeHealth(checks []metricsCheck) string {
summary := "ok"
for _, check := range checks {
switch check.Status {
case "failed":
return "failed"
case "warn":
summary = "warn"
}
}
return summary
}
func round1(value float64) float64 {
return float64(int(value*10+0.5)) / 10
}
+179
View File
@@ -0,0 +1,179 @@
package server
import (
"net/http"
"testing"
"time"
"glchat/internal/httpx"
)
// Тесты дашборда инстанса: доступ только администратору, состав метрик и
// healthcheck, отдельный лимит частоты для опроса раз в секунду (AGENT.md 7.19).
type metricsBody struct {
Metrics struct {
CPU struct {
Percent float64 `json:"percent"`
Cores float64 `json:"cores"`
Source string `json:"source"`
} `json:"cpu"`
Memory struct {
UsedBytes int64 `json:"used_bytes"`
TotalBytes int64 `json:"total_bytes"`
Percent float64 `json:"percent"`
Source string `json:"source"`
} `json:"memory"`
Disk struct {
UsedBytes int64 `json:"used_bytes"`
TotalBytes int64 `json:"total_bytes"`
FreeBytes int64 `json:"free_bytes"`
} `json:"disk"`
Database struct {
Bytes int64 `json:"bytes"`
WalBytes int64 `json:"wal_bytes"`
} `json:"database"`
Checks []struct {
Name string `json:"name"`
Status string `json:"status"`
Message string `json:"message"`
} `json:"checks"`
Health string `json:"health"`
UptimeSeconds int64 `json:"uptime_seconds"`
Version string `json:"version"`
Commit string `json:"commit"`
CollectedAt string `json:"collected_at"`
} `json:"metrics"`
}
// metricsServer готовит сервер с настоящими источниками метрик (cgroup/proc
// текущей машины) — тесты проверяют и разбор, и отдачу значений.
func metricsServer(t *testing.T) (*Server, *http.Cookie) {
t.Helper()
srv, _ := newTestServer(t)
cookie := registerAndLogin(t, srv, "metrics_admin", "metrics@example.com")
promoteAdmin(t, srv, "metrics@example.com")
return srv, cookie
}
func TestInstanceMetricsRequireAdmin(t *testing.T) {
srv, _ := newTestServer(t)
member := registerAndLogin(t, srv, "metrics_member", "metrics-member@example.com")
rec := doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", member)
if rec.Code != http.StatusForbidden {
t.Fatalf("обычный участник: статус %d, ожидался 403", rec.Code)
}
anonymous := doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "")
if anonymous.Code != http.StatusUnauthorized {
t.Fatalf("без сессии: статус %d, ожидался 401", anonymous.Code)
}
}
func TestInstanceMetricsPayload(t *testing.T) {
srv, cookie := metricsServer(t)
// Первый вызов только снимает показания процессора, второй считает дельту.
doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", cookie)
time.Sleep(20 * time.Millisecond)
rec := doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", cookie)
if rec.Code != http.StatusOK {
t.Fatalf("статус = %d, тело = %s", rec.Code, rec.Body.String())
}
payload := decodeResponse[metricsBody](t, rec).Metrics
if payload.Version == "" || payload.Commit != "deadbee" {
t.Fatalf("версия/коммит = %q/%q", payload.Version, payload.Commit)
}
if payload.UptimeSeconds < 0 || payload.CollectedAt == "" {
t.Fatalf("время работы/сбора = %d/%q", payload.UptimeSeconds, payload.CollectedAt)
}
// На macOS нет ни cgroup, ни /proc: источник недоступен, и ручка обязана
// деградировать мягко (проверка «metrics» в статусе warn).
switch payload.CPU.Source {
case "":
t.Log("источники процессора недоступны (не Linux) — проверка пропущена")
case "cgroup", "host":
if payload.CPU.Percent < 0 || payload.CPU.Percent > 100 || payload.CPU.Cores <= 0 {
t.Fatalf("процессор = %+v", payload.CPU)
}
default:
t.Fatalf("неизвестный источник процессора = %q", payload.CPU.Source)
}
if payload.Memory.TotalBytes == 0 {
t.Log("память недоступна (не Linux) — проверка пропущена")
} else if payload.Memory.UsedBytes <= 0 || payload.Memory.Percent <= 0 || payload.Memory.Percent > 100 {
t.Fatalf("память = %+v", payload.Memory)
}
if payload.Disk.TotalBytes <= 0 || payload.Disk.UsedBytes <= 0 || payload.Disk.FreeBytes <= 0 {
t.Fatalf("диск = %+v", payload.Disk)
}
if payload.Database.Bytes <= 0 {
t.Fatalf("размер базы = %d", payload.Database.Bytes)
}
// Healthcheck: приложение, база, диск, запись, gateway.
statuses := map[string]string{}
for _, check := range payload.Checks {
statuses[check.Name] = check.Status
}
for _, name := range []string{"app", "database", "disk", "storage", "gateway"} {
if statuses[name] != "ok" {
t.Fatalf("проверка %s = %q (все: %+v)", name, statuses[name], payload.Checks)
}
}
// На Linux без cgroup/proc не бывает, поэтому «metrics» появляется только
// на не-Linux: тогда health = warn, иначе ok.
if statuses["metrics"] == "warn" {
if payload.Health != "warn" {
t.Fatalf("состояние = %q при проблеме с метриками", payload.Health)
}
} else if payload.Health != "ok" {
t.Fatalf("итоговое состояние = %q, проверки: %+v", payload.Health, payload.Checks)
}
}
func TestInstanceMetricsRateLimitIsSeparate(t *testing.T) {
srv, cookie := metricsServer(t)
// Лимит дашборда — 120/мин с запасом 30: серия из 40 запросов проходит,
// хотя общий лимит API (120/мин) в тестах уже израсходован регистрацией.
for i := 0; i < 25; i++ {
rec := doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", cookie)
if rec.Code != http.StatusOK {
t.Fatalf("запрос %d: статус %d (%s)", i+1, rec.Code, rec.Body.String())
}
}
// Исчерпание своего лимита даёт 429 с подсказкой по паузе: ставим лимитер
// на один запрос и расходуем токен.
srv.metricsLimiter = httpx.NewRateLimiterWindow(1, time.Hour, 1)
doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", cookie)
rec := doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", cookie)
if rec.Code != http.StatusTooManyRequests {
t.Fatalf("статус = %d, ожидался 429", rec.Code)
}
payload := decodeResponse[struct {
Error struct {
Code string `json:"code"`
RetryAfterMS int `json:"retry_after_ms"`
} `json:"error"`
}](t, rec)
if payload.Error.Code != "rate_limited" || payload.Error.RetryAfterMS <= 0 {
t.Fatalf("ошибка = %+v", payload.Error)
}
}
func TestInstanceMetricsHealthCacheReusesChecks(t *testing.T) {
srv, cookie := metricsServer(t)
first := decodeResponse[metricsBody](t, doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", cookie)).Metrics
second := decodeResponse[metricsBody](t, doJSON(t, srv, http.MethodGet, "/api/v1/instance/metrics", "", cookie)).Metrics
if len(first.Checks) != len(second.Checks) {
t.Fatalf("набор проверок изменился: %d → %d", len(first.Checks), len(second.Checks))
}
// Проверки берутся из кэша (TTL 5 секунд) — время сбора метрик при этом
// обновляется на каждом запросе.
if first.CollectedAt == second.CollectedAt {
t.Fatal("метрики должны пересчитываться на каждый запрос")
}
}
+15 -1
View File
@@ -22,6 +22,7 @@ import (
"glchat/internal/permissions" "glchat/internal/permissions"
"glchat/internal/source" "glchat/internal/source"
"glchat/internal/store" "glchat/internal/store"
"glchat/internal/sysinfo"
"glchat/internal/voice" "glchat/internal/voice"
) )
@@ -66,6 +67,13 @@ type Server struct {
voice *voice.TokenIssuer voice *voice.TokenIssuer
// voiceAdmin — RoomService для модерации на стороне SFU. // voiceAdmin — RoomService для модерации на стороне SFU.
voiceAdmin *voice.AdminClient voiceAdmin *voice.AdminClient
// sysinfoPaths и cpuSampler — источники живых метрик дашборда инстанса.
sysinfoPaths sysinfo.Paths
cpuSampler *sysinfo.Sampler
// metricsLimiter — отдельный лимит для ручки метрик (опрос раз в секунду).
metricsLimiter *httpx.RateLimiter
// health — кэш тяжёлых проверок инстанса.
health healthCache
// presence — время последнего обновления last_seen по пользователю. // presence — время последнего обновления last_seen по пользователю.
presenceMu sync.Mutex presenceMu sync.Mutex
presence map[uint64]time.Time presence map[uint64]time.Time
@@ -101,9 +109,13 @@ func New(cfg config.Config, db *database.DB, logger *slog.Logger, deps Deps) *Se
slowmode: map[string]time.Time{}, slowmode: map[string]time.Time{},
presence: map[uint64]time.Time{}, presence: map[uint64]time.Time{},
webhookSeen: map[string]time.Time{}, webhookSeen: map[string]time.Time{},
// Дашборд опрашивает метрики раз в секунду: 120/мин с запасом.
metricsLimiter: httpx.NewRateLimiter(120, 30),
sysinfoPaths: sysinfo.DefaultPaths(),
voice: voice.NewIssuer(cfg.LiveKitAPIKey, cfg.LiveKitAPISecret, cfg.LiveKitTokenTTL), voice: voice.NewIssuer(cfg.LiveKitAPIKey, cfg.LiveKitAPISecret, cfg.LiveKitTokenTTL),
voiceAdmin: voice.NewAdminClient(voice.NewIssuer(cfg.LiveKitAPIKey, cfg.LiveKitAPISecret, cfg.LiveKitTokenTTL), cfg.LiveKitAPIURL), voiceAdmin: voice.NewAdminClient(voice.NewIssuer(cfg.LiveKitAPIKey, cfg.LiveKitAPISecret, cfg.LiveKitTokenTTL), cfg.LiveKitAPIURL),
} }
s.cpuSampler = sysinfo.NewSampler(s.sysinfoPaths)
switch { switch {
case deps.Permissions != nil: case deps.Permissions != nil:
s.perms = deps.Permissions s.perms = deps.Permissions
@@ -115,7 +127,8 @@ func New(cfg config.Config, db *database.DB, logger *slog.Logger, deps Deps) *Se
router.Route("/api/v1", func(apiRouter chi.Router) { router.Route("/api/v1", func(apiRouter chi.Router) {
// Сессия резолвится один раз на запрос: huma-ручки читают её из контекста. // Сессия резолвится один раз на запрос: huma-ручки читают её из контекста.
apiRouter.Use(s.sessionContext) apiRouter.Use(s.sessionContext)
apiRouter.Use(s.apiLimiter.Middleware(apiRateLimitKey)) // Метрики дашборда вне общего лимита: у них свой (см. api_metrics.go).
apiRouter.Use(s.apiLimiter.MiddlewareExcept(apiRateLimitKey, metricsPath))
s.api = s.registerAPI(apiRouter) s.api = s.registerAPI(apiRouter)
s.registerMetaRoutes(s.api) s.registerMetaRoutes(s.api)
s.registerAuthRoutes(apiRouter) s.registerAuthRoutes(apiRouter)
@@ -123,6 +136,7 @@ func New(cfg config.Config, db *database.DB, logger *slog.Logger, deps Deps) *Se
s.registerUserRoutes(s.api) s.registerUserRoutes(s.api)
s.registerGuildRoutes(s.api) s.registerGuildRoutes(s.api)
s.registerInstanceRoutes(s.api) s.registerInstanceRoutes(s.api)
s.registerMetricsRoutes(s.api)
s.registerMessageRoutes(s.api) s.registerMessageRoutes(s.api)
s.registerInviteRoutes(s.api) s.registerInviteRoutes(s.api)
s.registerFileRoutes(s.api, apiRouter) s.registerFileRoutes(s.api, apiRouter)
+378
View File
@@ -0,0 +1,378 @@
// Package sysinfo читает нагрузку хоста и контейнера для дашборда инстанса:
// процессор, память, диск и размеры файлов. Значения берутся из cgroup v2
// (контейнер ограничен профилем установки), а если лимитов нет — из /proc.
package sysinfo
import (
"bufio"
"errors"
"fmt"
"math"
"os"
"path/filepath"
"runtime"
"strconv"
"strings"
"sync"
"syscall"
"time"
)
// userHZ — тиков в секунду для /proc/stat: в Linux это всегда 100, независимо
// от CONFIG_HZ ядра.
const userHZ = 100
// Paths — пути к источникам данных. Вынесены в структуру, чтобы тесты
// подменяли их временными файлами, а не зависели от окружения.
type Paths struct {
ProcStat string
MemInfo string
CgroupDir string
WallClock func() time.Time
}
// DefaultPaths возвращает источники по умолчанию для Linux-контейнера.
func DefaultPaths() Paths {
return Paths{
ProcStat: "/proc/stat",
MemInfo: "/proc/meminfo",
CgroupDir: "/sys/fs/cgroup",
WallClock: time.Now,
}
}
// CPUSample — снимок процессорного времени.
type CPUSample struct {
// UsageMicros — накопленное использованное время.
UsageMicros float64
// TotalMicros — накопленное доступное время (режим хоста); в режиме
// cgroup не используется, там ёмкость считается по квоте и времени.
TotalMicros float64
// QuotaPerSecond — доступное время в микросекундах за секунду (cgroup).
QuotaPerSecond float64
At time.Time
}
// Sampler считает загрузку процессора между вызовами: вызовы раз в секунду
// дают нагрузку за последнюю секунду, редкие вызовы — среднее за паузу.
type Sampler struct {
paths Paths
mu sync.Mutex
last *CPUSample
}
// NewSampler создаёт сэмплер процессорного времени.
func NewSampler(paths Paths) *Sampler {
return &Sampler{paths: paths}
}
// CPUPercent возвращает загрузку в процентах от доступного времени и число
// ядер. Первый вызов без предыдущего снимка возвращает 0: дельты ещё нет.
func (s *Sampler) CPUPercent() (percent float64, cores float64, source string, err error) {
sample, cores, source, err := s.readCPU()
if err != nil {
return 0, cores, source, err
}
s.mu.Lock()
previous := s.last
s.last = &sample
s.mu.Unlock()
if previous == nil {
return 0, cores, source, nil
}
used := sample.UsageMicros - previous.UsageMicros
if used < 0 {
used = 0
}
elapsed := sample.At.Sub(previous.At).Seconds()
if elapsed <= 0 {
return 0, cores, source, nil
}
var capacity float64
if sample.QuotaPerSecond > 0 {
// Контейнер с лимитом: доступное время = квота × прошедшее время.
capacity = sample.QuotaPerSecond * elapsed
} else {
// Вся машина: доступное время = прирост суммарных тиков всех ядер.
capacity = sample.TotalMicros - previous.TotalMicros
}
if capacity <= 0 {
return 0, cores, source, nil
}
return used / capacity * 100, cores, source, nil
}
// readCPU предпочитает cgroup v2: он показывает нагрузку контейнера и его
// лимит, а /proc/stat — всю машину целиком.
func (s *Sampler) readCPU() (CPUSample, float64, string, error) {
if sample, cores, ok := s.readCgroupCPU(); ok {
return sample, cores, "cgroup", nil
}
sample, cores, err := s.readProcStat()
if err != nil {
return CPUSample{}, 0, "", err
}
return sample, cores, "host", nil
}
// readCgroupCPU читает cpu.stat и cpu.max (cgroup v2). Квота задаётся парой
// «quota period»; значение «max» означает отсутствие лимита.
func (s *Sampler) readCgroupCPU() (CPUSample, float64, bool) {
usage, err := readMicros(filepath.Join(s.paths.CgroupDir, "cpu.stat"), "usage_usec")
if err != nil {
return CPUSample{}, 0, false
}
cores, ok := s.cpuQuota()
if !ok {
return CPUSample{}, 0, false
}
return CPUSample{
UsageMicros: usage,
QuotaPerSecond: cores * 1e6,
At: s.now(),
}, cores, true
}
// cpuQuota возвращает число доступных ядер по cpu.max.
func (s *Sampler) cpuQuota() (float64, bool) {
raw, err := os.ReadFile(filepath.Join(s.paths.CgroupDir, "cpu.max"))
if err != nil {
return 0, false
}
fields := strings.Fields(string(raw))
if len(fields) != 2 || fields[0] == "max" {
return 0, false
}
quota, err := strconv.ParseFloat(fields[0], 64)
if err != nil {
return 0, false
}
period, err := strconv.ParseFloat(fields[1], 64)
if err != nil || period <= 0 {
return 0, false
}
cores := quota / period
if cores <= 0 {
return 0, false
}
return cores, true
}
// readProcStat считает загрузку всей машины по первой строке /proc/stat.
func (s *Sampler) readProcStat() (CPUSample, float64, error) {
file, err := os.Open(s.paths.ProcStat)
if err != nil {
return CPUSample{}, 0, fmt.Errorf("sysinfo.cpu: %w", err)
}
defer func() { _ = file.Close() }()
scanner := bufio.NewScanner(file)
for scanner.Scan() {
fields := strings.Fields(scanner.Text())
if len(fields) < 5 || fields[0] != "cpu" {
continue
}
var total, idle float64
for index, field := range fields[1:] {
value, err := strconv.ParseFloat(field, 64)
if err != nil {
return CPUSample{}, 0, fmt.Errorf("sysinfo.cpu: %w", err)
}
total += value
// idle и iowait — время, когда процессор не работал.
if index == 3 || index == 4 {
idle += value
}
}
if total <= 0 {
return CPUSample{}, 0, errors.New("sysinfo.cpu: пустой счётчик")
}
return CPUSample{
UsageMicros: ticksToMicros(total - idle),
TotalMicros: ticksToMicros(total),
At: s.now(),
}, float64(runtime.NumCPU()), nil
}
if err := scanner.Err(); err != nil {
return CPUSample{}, 0, fmt.Errorf("sysinfo.cpu: %w", err)
}
return CPUSample{}, 0, errors.New("sysinfo.cpu: строка cpu не найдена")
}
func (s *Sampler) now() time.Time {
if s.paths.WallClock == nil {
return time.Now()
}
return s.paths.WallClock()
}
// Memory — занятая и доступная память в байтах.
type Memory struct {
UsedBytes int64
TotalBytes int64
Source string
}
// Percent возвращает заполнение памяти в процентах.
func (m Memory) Percent() float64 {
if m.TotalBytes <= 0 {
return 0
}
return float64(m.UsedBytes) / float64(m.TotalBytes) * 100
}
// ReadMemory отдаёт лимит и потребление контейнера (cgroup v2), а при его
// отсутствии — память всей машины по /proc/meminfo.
func ReadMemory(paths Paths) (Memory, error) {
if memory, ok := readCgroupMemory(paths.CgroupDir); ok {
return memory, nil
}
return readMemInfo(paths.MemInfo)
}
func readCgroupMemory(dir string) (Memory, bool) {
current, err := readInt(filepath.Join(dir, "memory.current"))
if err != nil {
return Memory{}, false
}
total, err := readInt(filepath.Join(dir, "memory.max"))
if err != nil || total <= 0 {
// memory.max = "max" или файла нет: лимит не выставлен.
return Memory{}, false
}
return Memory{UsedBytes: current, TotalBytes: total, Source: "cgroup"}, true
}
func readMemInfo(path string) (Memory, error) {
file, err := os.Open(path) //nolint:gosec // путь задаёт сам инстанс (по умолчанию /proc/meminfo)
if err != nil {
return Memory{}, fmt.Errorf("sysinfo.memory: %w", err)
}
defer func() { _ = file.Close() }()
var total, available int64
scanner := bufio.NewScanner(file)
for scanner.Scan() {
fields := strings.Fields(scanner.Text())
if len(fields) < 2 {
continue
}
value, err := strconv.ParseInt(fields[1], 10, 64)
if err != nil {
continue
}
switch strings.TrimSuffix(fields[0], ":") {
case "MemTotal":
total = value * 1024
case "MemAvailable":
available = value * 1024
}
}
if err := scanner.Err(); err != nil {
return Memory{}, fmt.Errorf("sysinfo.memory: %w", err)
}
if total <= 0 {
return Memory{}, errors.New("sysinfo.memory: MemTotal не найден")
}
return Memory{UsedBytes: total - available, TotalBytes: total, Source: "host"}, nil
}
// Disk — занятое и общее место раздела, на котором лежат данные инстанса.
type Disk struct {
UsedBytes int64
TotalBytes int64
FreeBytes int64
}
// ReadDisk считает место по пути данных (в контейнере это примонтированный
// каталог хоста).
func ReadDisk(path string) (Disk, error) {
var stats syscall.Statfs_t
if err := syscall.Statfs(path, &stats); err != nil {
return Disk{}, fmt.Errorf("sysinfo.disk: %w", err)
}
blockSize := toInt64From32(stats.Bsize)
total := saturatingMul(toInt64(stats.Blocks), blockSize)
free := saturatingMul(toInt64(stats.Bavail), blockSize)
return Disk{UsedBytes: total - free, TotalBytes: total, FreeBytes: free}, nil
}
// FileBytes суммирует размеры существующих файлов: отсутствующие считаются
// нулём, поэтому WAL и SHM можно передавать всегда.
func FileBytes(paths ...string) int64 {
var total int64
for _, path := range paths {
info, err := os.Stat(path)
if err != nil || info.IsDir() {
continue
}
total += info.Size()
}
return total
}
// toInt64 переводит беззнаковый размер в int64 с насыщением: значения statfs
// для реальных разделов далеки от предела, но переполнение недопустимо.
func toInt64(value uint64) int64 {
if value > math.MaxInt64 {
return math.MaxInt64
}
return int64(value)
}
// toInt64From32 переводит 32-битное беззнаковое значение (размер блока statfs).
func toInt64From32(value uint32) int64 {
return int64(value)
}
// saturatingMul умножает с насыщением: размер раздела не должен переполняться.
func saturatingMul(a, b int64) int64 {
if a <= 0 || b <= 0 {
return 0
}
if a > math.MaxInt64/b {
return math.MaxInt64
}
return a * b
}
func ticksToMicros(ticks float64) float64 {
return ticks / userHZ * 1e6
}
func readInt(path string) (int64, error) {
raw, err := os.ReadFile(path) //nolint:gosec // путь задаёт сам инстанс (cgroup/proc или временный каталог в тестах)
if err != nil {
return 0, err
}
return strconv.ParseInt(strings.TrimSpace(string(raw)), 10, 64)
}
func readMicros(path, key string) (float64, error) {
file, err := os.Open(path) //nolint:gosec // путь задаёт сам инстанс (cgroup/proc)
if err != nil {
return 0, err
}
defer func() { _ = file.Close() }()
scanner := bufio.NewScanner(file)
for scanner.Scan() {
fields := strings.Fields(scanner.Text())
if len(fields) != 2 || fields[0] != key {
continue
}
value, err := strconv.ParseFloat(fields[1], 64)
if err != nil {
return 0, err
}
return value, nil
}
if err := scanner.Err(); err != nil {
return 0, err
}
return 0, fmt.Errorf("sysinfo: %s не найден в %s", key, path)
}
+177
View File
@@ -0,0 +1,177 @@
package sysinfo_test
import (
"os"
"path/filepath"
"testing"
"time"
"glchat/internal/sysinfo"
)
// Тесты работают на временных файлах: значения cgroup и /proc подменяются,
// поэтому проверки не зависят от окружения и работают в контейнере CI.
func writeFile(t *testing.T, dir, name, content string) string {
t.Helper()
path := filepath.Join(dir, name)
if err := os.WriteFile(path, []byte(content), 0o600); err != nil {
t.Fatalf("запись %s: %v", name, err)
}
return path
}
func TestCPUPercentUsesCgroupQuota(t *testing.T) {
dir := t.TempDir()
writeFile(t, dir, "cpu.max", "200000 100000\n") // 2 ядра
writeFile(t, dir, "cpu.stat", "usage_usec 1000000\nuser_usec 900000\nsystem_usec 100000\n")
now := time.Unix(1_700_000_000, 0)
paths := sysinfo.Paths{
ProcStat: filepath.Join(dir, "stat"),
MemInfo: filepath.Join(dir, "meminfo"),
CgroupDir: dir,
WallClock: func() time.Time { return now },
}
sampler := sysinfo.NewSampler(paths)
// Первый вызов только запоминает снимок.
if percent, _, source, err := sampler.CPUPercent(); err != nil || percent != 0 || source != "cgroup" {
t.Fatalf("первый замер = %v, %s, %v", percent, source, err)
}
// За секунду израсходована половина квоты (1 с из 2 ядро-секунд) → 50 %.
now = now.Add(time.Second)
writeFile(t, dir, "cpu.stat", "usage_usec 2000000\n")
percent, cores, source, err := sampler.CPUPercent()
if err != nil {
t.Fatalf("второй замер: %v", err)
}
if cores != 2 || source != "cgroup" {
t.Fatalf("ядер=%v источник=%s", cores, source)
}
if percent < 49 || percent > 51 {
t.Fatalf("загрузка = %.2f %%, ожидалось ~50 %%", percent)
}
}
func TestCPUPercentFallsBackToProcStat(t *testing.T) {
dir := t.TempDir()
statPath := writeFile(t, dir, "stat", "cpu 100 0 100 800 0 0 0 0 0 0\n")
writeFile(t, dir, "meminfo", "MemTotal: 1000 kB\nMemAvailable: 400 kB\n")
now := time.Unix(1_700_000_000, 0)
sampler := sysinfo.NewSampler(sysinfo.Paths{
ProcStat: statPath,
MemInfo: filepath.Join(dir, "meminfo"),
CgroupDir: filepath.Join(dir, "нет-cgroup"),
WallClock: func() time.Time { return now },
})
if _, _, source, err := sampler.CPUPercent(); err != nil || source != "host" {
t.Fatalf("источник = %s, %v", source, err)
}
// За интервал прибавилось 150 тиков занятости и 150 простоя → 50 %.
now = now.Add(time.Second)
writeFile(t, dir, "stat", "cpu 175 0 175 950 0 0 0 0 0 0\n")
percent, cores, source, err := sampler.CPUPercent()
if err != nil {
t.Fatalf("второй замер: %v", err)
}
if source != "host" || cores <= 0 {
t.Fatalf("источник=%s ядер=%v", source, cores)
}
if percent < 49 || percent > 51 {
t.Fatalf("загрузка = %.2f %%, ожидалось ~50 %%", percent)
}
}
func TestReadMemoryPrefersCgroupLimit(t *testing.T) {
dir := t.TempDir()
writeFile(t, dir, "memory.current", "15728640\n")
writeFile(t, dir, "memory.max", "805306368\n")
writeFile(t, dir, "meminfo", "MemTotal: 3479592 kB\nMemAvailable: 1000000 kB\n")
memory, err := sysinfo.ReadMemory(sysinfo.Paths{
MemInfo: filepath.Join(dir, "meminfo"),
CgroupDir: dir,
ProcStat: filepath.Join(dir, "stat"),
})
if err != nil {
t.Fatalf("память: %v", err)
}
if memory.Source != "cgroup" || memory.TotalBytes != 805306368 || memory.UsedBytes != 15728640 {
t.Fatalf("память = %+v", memory)
}
if percent := memory.Percent(); percent < 1.9 || percent > 2.0 {
t.Fatalf("заполнение = %.2f %%", percent)
}
}
func TestReadMemoryFallsBackToMemInfo(t *testing.T) {
dir := t.TempDir()
writeFile(t, dir, "meminfo", "MemTotal: 1000 kB\nMemAvailable: 250 kB\n")
memory, err := sysinfo.ReadMemory(sysinfo.Paths{
MemInfo: filepath.Join(dir, "meminfo"),
CgroupDir: filepath.Join(dir, "нет-cgroup"),
ProcStat: filepath.Join(dir, "stat"),
})
if err != nil {
t.Fatalf("память: %v", err)
}
if memory.Source != "host" || memory.TotalBytes != 1024_000 || memory.UsedBytes != 768_000 {
t.Fatalf("память = %+v", memory)
}
}
func TestReadMemoryUnlimitedCgroupMax(t *testing.T) {
dir := t.TempDir()
writeFile(t, dir, "memory.current", "1000\n")
writeFile(t, dir, "memory.max", "max\n")
writeFile(t, dir, "meminfo", "MemTotal: 2000 kB\nMemAvailable: 500 kB\n")
memory, err := sysinfo.ReadMemory(sysinfo.Paths{
MemInfo: filepath.Join(dir, "meminfo"),
CgroupDir: dir,
ProcStat: filepath.Join(dir, "stat"),
})
if err != nil {
t.Fatalf("память: %v", err)
}
if memory.Source != "host" {
t.Fatalf("при memory.max=max ожидался хостовый режим, получено %+v", memory)
}
}
func TestReadDiskAndFileBytes(t *testing.T) {
dir := t.TempDir()
disk, err := sysinfo.ReadDisk(dir)
if err != nil {
t.Fatalf("диск: %v", err)
}
if disk.TotalBytes <= 0 || disk.UsedBytes <= 0 || disk.FreeBytes <= 0 {
t.Fatalf("диск = %+v", disk)
}
if disk.UsedBytes+disk.FreeBytes > disk.TotalBytes {
t.Fatalf("занято+свободно больше общего: %+v", disk)
}
writeFile(t, dir, "glchat.db", "0123456789")
writeFile(t, dir, "glchat.db-wal", "01234")
if got := sysinfo.FileBytes(filepath.Join(dir, "glchat.db"), filepath.Join(dir, "glchat.db-wal"), filepath.Join(dir, "нет.db")); got != 15 {
t.Fatalf("размер файлов = %d, ожидалось 15", got)
}
}
func TestCPUPercentWithoutProcStat(t *testing.T) {
sampler := sysinfo.NewSampler(sysinfo.Paths{
ProcStat: filepath.Join(t.TempDir(), "нет"),
MemInfo: filepath.Join(t.TempDir(), "нет"),
CgroupDir: filepath.Join(t.TempDir(), "нет"),
})
if _, _, _, err := sampler.CPUPercent(); err == nil {
t.Fatal("без источников ожидалась ошибка")
}
}
+20 -2
View File
@@ -110,17 +110,35 @@ func (c *AdminClient) ListParticipants(ctx context.Context, room string) ([]Part
return response.Participants, nil return response.Participants, nil
} }
// Ping проверяет доступность RoomService: healthcheck дашборда заодно
// убеждается, что ключи LiveKit валидны (их неверность ломает модерацию).
func (c *AdminClient) Ping(ctx context.Context) error {
if !c.Enabled() {
return ErrDisabled
}
token, err := c.issuer.IssueRoomList()
if err != nil {
return err
}
return c.callWithToken(ctx, token, "ListRooms", map[string]any{}, nil)
}
// call выполняет запрос к Twirp-методу RoomService с административным токеном // call выполняет запрос к Twirp-методу RoomService с административным токеном
// этой комнаты: LiveKit принимает `roomAdmin` только вместе с именем комнаты. // этой комнаты: LiveKit принимает `roomAdmin` только вместе с именем комнаты.
func (c *AdminClient) call(ctx context.Context, room, method string, body any, out any) error { func (c *AdminClient) call(ctx context.Context, room, method string, body any, out any) error {
if !c.Enabled() { if !c.Enabled() {
return ErrDisabled return ErrDisabled
} }
payload, err := json.Marshal(body) token, err := c.issuer.IssueAdmin(room)
if err != nil { if err != nil {
return err return err
} }
token, err := c.issuer.IssueAdmin(room) return c.callWithToken(ctx, token, method, body, out)
}
// callWithToken выполняет запрос к Twirp-методу с готовым токеном.
func (c *AdminClient) callWithToken(ctx context.Context, token, method string, body any, out any) error {
payload, err := json.Marshal(body)
if err != nil { if err != nil {
return err return err
} }
+23 -14
View File
@@ -87,6 +87,11 @@ func (i *TokenIssuer) Issue(room, identity, displayName string, grants Grants) (
claims.NotBefore = now.Add(-10 * time.Second).Unix() claims.NotBefore = now.Add(-10 * time.Second).Unix()
claims.ExpiresAt = now.Add(i.ttl).Unix() claims.ExpiresAt = now.Add(i.ttl).Unix()
return i.sign(claims)
}
// sign подписывает произвольные claims алгоритмом HS256.
func (i *TokenIssuer) sign(claims any) (string, error) {
header, err := json.Marshal(map[string]string{"alg": "HS256", "typ": "JWT"}) header, err := json.Marshal(map[string]string{"alg": "HS256", "typ": "JWT"})
if err != nil { if err != nil {
return "", err return "", err
@@ -142,20 +147,7 @@ func (i *TokenIssuer) IssueAdmin(room string) (string, error) {
"roomAdmin": true, "roomAdmin": true,
}, },
} }
header, err := json.Marshal(map[string]string{"alg": "HS256", "typ": "JWT"}) return i.sign(claims)
if err != nil {
return "", err
}
payload, err := json.Marshal(claims)
if err != nil {
return "", err
}
signingInput := encodeSegment(header) + "." + encodeSegment(payload)
mac := hmac.New(sha256.New, []byte(i.apiSecret))
if _, err := mac.Write([]byte(signingInput)); err != nil {
return "", err
}
return signingInput + "." + base64.RawURLEncoding.EncodeToString(mac.Sum(nil)), nil
} }
// VerifyWebhook проверяет подпись вебхука LiveKit: это JWT, подписанный тем же // VerifyWebhook проверяет подпись вебхука LiveKit: это JWT, подписанный тем же
@@ -272,6 +264,23 @@ func ResolveLiveKitURL(raw string) (string, error) {
return strings.TrimSuffix(parsed.String(), "/"), nil return strings.TrimSuffix(parsed.String(), "/"), nil
} }
// IssueRoomList подписывает служебный токен для запросов, которым не нужна
// конкретная комната (проверка доступности RoomService в healthcheck).
func (i *TokenIssuer) IssueRoomList() (string, error) {
if !i.Enabled() {
return "", ErrDisabled
}
return i.sign(map[string]any{
"iss": i.apiKey,
"sub": i.apiKey,
"nbf": i.now().UTC().Add(-10 * time.Second).Unix(),
"exp": i.now().UTC().Add(5 * time.Minute).Unix(),
"video": map[string]any{
"roomList": true,
},
})
}
// VoiceURLForClient отдаёт адрес LiveKit так, как его должен использовать // VoiceURLForClient отдаёт адрес LiveKit так, как его должен использовать
// клиент: пусто, если голос не настроен. // клиент: пусто, если голос не настроен.
func (i *TokenIssuer) VoiceURLForClient(publicURL string) string { func (i *TokenIssuer) VoiceURLForClient(publicURL string) string {