Files
glchat/cmd/loadgen/main.go
T
grendervill 66dd07eb7b feat(ops): генератор нагрузки и подготовка аккаунтов для теста
- cmd/loadgen: держит заданное число WS-клиентов Gateway, шлёт сообщения с
  нужной частотой, измеряет задержку доставки (p50/p95/p99/max), ответы API,
  разрывы соединений и пишет отчёт в JSON;
- glchat create-user (CLI + обёртка): создаёт аккаунты в обход выключенной
  регистрации, пакетно (--count/--prefix), при --sessions сразу выдаёт сессии
  и складывает учётные данные в файл 600 — иначе сотня клиентов не сможет
  войти из-за лимита попыток по IP;
- auth: CreateUserByOperator и IssueSessionForOperator с тестом;
- web/e2e/voice-ten.spec.ts: десять участников в одной голосовой комнате,
  проверка плиток, входящего аудио и слоя 1080p60 у клиента; сессии берутся
  из файла (без формы входа);
- .gitignore: каталоги вывода Playwright test-results*.
2026-09-22 20:00:33 +03:00

578 lines
18 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// loadgen — генератор нагрузки для проверки целевых показателей AGENT.md 9.1,
// 11.4: держит заданное число WebSocket-клиентов Gateway, шлёт сообщения с
// нужной частотой и измеряет задержку доставки и ответы API.
//
// Аккаунты готовит команда обслуживания на сервере:
//
// glchat create-user --count 100 --prefix loadtest --password '...' \
// --out /tmp/loadtest-users.json --sessions
//
// Затем генератор запускается с машины-генератора (не с сервера, иначе он
// отбирает у сервера процессор и портит замеры):
//
// go run ./cmd/loadgen -base https://gl.mhspx.su \
// -creds /tmp/loadtest-users.json -invite <код> \
// -clients 100 -rate 20 -duration 60s
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"math/rand"
"net/http"
"os"
"sort"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/coder/websocket"
)
type credentials struct {
ID string `json:"id"`
Username string `json:"username"`
Email string `json:"email"`
Password string `json:"password"`
Token string `json:"token"`
}
// envelope — кадр Gateway (AGENT.md 8.2).
type envelope struct {
Op int `json:"op"`
T string `json:"t,omitempty"`
D json.RawMessage `json:"d,omitempty"`
S int64 `json:"s,omitempty"`
}
const (
opDispatch = 0
opHeartbeat = 1
opIdentify = 2
opHello = 10
opHeartbeatAck = 11
)
// report — итог прогона: его печатают и складывают в JSON.
type report struct {
Clients int `json:"clients"`
Rate float64 `json:"messages_per_second"`
DurationSeconds float64 `json:"duration_seconds"`
Sent int64 `json:"messages_sent"`
SendErrors int64 `json:"send_errors"`
Received int64 `json:"events_received"`
Delivered int64 `json:"messages_delivered"`
DeliveryP50MS float64 `json:"delivery_p50_ms"`
DeliveryP95MS float64 `json:"delivery_p95_ms"`
DeliveryP99MS float64 `json:"delivery_p99_ms"`
DeliveryMaxMS float64 `json:"delivery_max_ms"`
APIP50MS float64 `json:"api_p50_ms"`
APIP95MS float64 `json:"api_p95_ms"`
APIMaxMS float64 `json:"api_max_ms"`
APIRequests int64 `json:"api_requests"`
APIErrors int64 `json:"api_errors"`
ReadyClients int `json:"ready_clients"`
Disconnects int64 `json:"disconnects"`
StartedAt time.Time `json:"started_at"`
}
func main() {
base := flag.String("base", "https://gl.mhspx.su", "адрес инстанса")
credsPath := flag.String("creds", "", "файл с учётными данными (glchat create-user --out --sessions)")
invite := flag.String("invite", "", "код приглашения в тестовый сервер (необязательно)")
channelID := flag.String("channel", "", "комната для сообщений (по умолчанию — первая текстовая сервера)")
guildID := flag.String("guild", "", "сервер, в котором работают клиенты")
clients := flag.Int("clients", 100, "число WS-клиентов")
rate := flag.Float64("rate", 20, "частота сообщений в секунду (все клиенты вместе)")
duration := flag.Duration("duration", 60*time.Second, "длительность нагрузки")
warmup := flag.Duration("warmup", 5*time.Second, "прогрев до начала замера")
apiEvery := flag.Duration("api-interval", 2*time.Second, "частота проверок API одним клиентом")
jsonOut := flag.String("json", "", "куда записать отчёт в JSON")
flag.Parse()
if *credsPath == "" {
fmt.Fprintln(os.Stderr, "loadgen: укажите -creds")
os.Exit(2)
}
all, err := readCredentials(*credsPath)
if err != nil {
fmt.Fprintf(os.Stderr, "loadgen: %v\n", err)
os.Exit(1)
}
if *clients > len(all) {
fmt.Fprintf(os.Stderr, "loadgen: в файле %d аккаунтов, нужно %d\n", len(all), *clients)
os.Exit(2)
}
baseURL := strings.TrimRight(*base, "/")
runner := &runner{
baseURL: baseURL,
clients: all[:*clients],
rate: *rate,
warmup: *warmup,
apiEvery: *apiEvery,
}
if *invite != "" {
if err := runner.joinGuild(context.Background(), *invite); err != nil {
fmt.Fprintf(os.Stderr, "loadgen: %v\n", err)
os.Exit(1)
}
}
target, err := runner.resolveChannel(context.Background(), *guildID, *channelID)
if err != nil {
fmt.Fprintf(os.Stderr, "loadgen: %v\n", err)
os.Exit(1)
}
runner.guildID = target.guildID
runner.channelID = target.channelID
fmt.Printf("сервер %s, комната %s, клиентов %d, темп %.1f msg/s, длительность %s\n",
target.guildID, target.channelID, len(runner.clients), *rate, duration.String())
final := runner.run(context.Background(), *duration)
printReport(final)
if *jsonOut != "" {
if err := writeReport(*jsonOut, final); err != nil {
fmt.Fprintf(os.Stderr, "loadgen: %v\n", err)
os.Exit(1)
}
fmt.Printf("отчёт: %s\n", *jsonOut)
}
}
func readCredentials(path string) ([]credentials, error) {
// Файл создаёт оператор командой create-user: путь задаётся явно.
raw, err := os.ReadFile(path) //nolint:gosec
if err != nil {
return nil, err
}
var list []credentials
if err := json.Unmarshal(raw, &list); err != nil {
return nil, fmt.Errorf("разобрать %s: %w", path, err)
}
for _, item := range list {
if item.Token == "" {
return nil, errors.New("в файле нет сессий: создайте аккаунты с флагом --sessions")
}
}
return list, nil
}
type channelTarget struct {
guildID string
channelID string
}
type runner struct {
baseURL string
clients []credentials
rate float64
warmup time.Duration
apiEvery time.Duration
client *http.Client
guildID string
channelID string
sent atomic.Int64
sendErrors atomic.Int64
received atomic.Int64
delivered atomic.Int64
disconnect atomic.Int64
apiCount atomic.Int64
apiErrors atomic.Int64
ready atomic.Int64
mu sync.Mutex
sendTimes map[string]time.Time
latencies []float64
apiLatency []float64
}
func (r *runner) httpClient() *http.Client {
if r.client == nil {
r.client = &http.Client{Timeout: 15 * time.Second}
}
return r.client
}
// joinGuild вступает в тестовый сервер по приглашению от имени всех клиентов.
func (r *runner) joinGuild(ctx context.Context, code string) error {
failed := atomic.Int64{}
var wg sync.WaitGroup
for _, item := range r.clients {
wg.Add(1)
go func(item credentials) {
defer wg.Done()
response, err := r.do(ctx, item.Token, http.MethodPost, "/api/v1/invites/"+code, nil)
if err != nil {
failed.Add(1)
return
}
defer func() { _ = response.Body.Close() }()
// 409 — уже участник: это не ошибка прогона.
if response.StatusCode >= 400 && response.StatusCode != http.StatusConflict {
failed.Add(1)
}
}(item)
}
wg.Wait()
if failed.Load() > 0 {
return fmt.Errorf("не все клиенты вошли в сервер: ошибок %d", failed.Load())
}
fmt.Printf("клиенты вошли в сервер по приглашению %s\n", code)
return nil
}
// resolveChannel выбирает сервер и текстовую комнату для нагрузки.
func (r *runner) resolveChannel(ctx context.Context, guildID, channelID string) (channelTarget, error) {
if guildID == "" {
response, err := r.do(ctx, r.clients[0].Token, http.MethodGet, "/api/v1/users/@me/guilds", nil)
if err != nil {
return channelTarget{}, err
}
defer func() { _ = response.Body.Close() }()
var payload struct {
Guilds []struct {
ID string `json:"id"`
} `json:"guilds"`
}
if err := json.NewDecoder(response.Body).Decode(&payload); err != nil {
return channelTarget{}, err
}
if len(payload.Guilds) == 0 {
return channelTarget{}, errors.New("у клиента нет серверов: передайте -guild и -invite")
}
guildID = payload.Guilds[0].ID
}
response, err := r.do(ctx, r.clients[0].Token, http.MethodGet, "/api/v1/guilds/"+guildID+"/channels", nil)
if err != nil {
return channelTarget{}, err
}
defer func() { _ = response.Body.Close() }()
var payload struct {
Channels []struct {
ID string `json:"id"`
Type string `json:"type"`
} `json:"channels"`
}
if err := json.NewDecoder(response.Body).Decode(&payload); err != nil {
return channelTarget{}, err
}
if channelID == "" {
for _, channel := range payload.Channels {
if channel.Type == "text" {
channelID = channel.ID
break
}
}
}
if channelID == "" {
return channelTarget{}, errors.New("в сервере нет текстовых комнат")
}
return channelTarget{guildID: guildID, channelID: channelID}, nil
}
// run поднимает клиентов, греет их, затем держит нагрузку заданное время.
func (r *runner) run(ctx context.Context, duration time.Duration) report {
r.sendTimes = make(map[string]time.Time, 4096)
started := time.Now().UTC()
runCtx, cancel := context.WithCancel(ctx)
defer cancel()
var wg sync.WaitGroup
for index := range r.clients {
wg.Add(1)
go func(index int) {
defer wg.Done()
r.clientLoop(runCtx, index)
}(index)
}
// Ждём READY у большинства клиентов: иначе замер пойдёт по «холодным» сессиям.
deadline := time.Now().Add(30 * time.Second)
for time.Now().Before(deadline) && r.ready.Load() < int64(len(r.clients)*9/10) {
time.Sleep(200 * time.Millisecond)
}
fmt.Printf("готовы к работе: %d из %d клиентов\n", r.ready.Load(), len(r.clients))
time.Sleep(r.warmup)
sendCtx, stopSending := context.WithTimeout(runCtx, duration)
defer stopSending()
r.sendLoop(sendCtx)
// Даём последним сообщениям дойти до клиентов.
time.Sleep(2 * time.Second)
cancel()
wg.Wait()
final := r.report(started, duration)
return final
}
// sendLoop шлёт сообщения с заданной частотой, распределяя их по клиентам.
func (r *runner) sendLoop(ctx context.Context) {
interval := time.Duration(float64(time.Second) / r.rate)
ticker := time.NewTicker(interval)
defer ticker.Stop()
random := rand.New(rand.NewSource(time.Now().UnixNano()))
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
sender := r.clients[random.Intn(len(r.clients))]
go r.sendMessage(ctx, sender)
}
}
}
// sendMessage отправляет сообщение и запоминает время отправки по id.
func (r *runner) sendMessage(ctx context.Context, sender credentials) {
body := map[string]any{"content": fmt.Sprintf("нагрузка %d", time.Now().UnixNano())}
response, err := r.do(ctx, sender.Token, http.MethodPost, "/api/v1/channels/"+r.channelID+"/messages", body)
if err != nil {
r.sendErrors.Add(1)
return
}
defer func() { _ = response.Body.Close() }()
if response.StatusCode != http.StatusOK {
r.sendErrors.Add(1)
return
}
var payload struct {
Message struct {
ID string `json:"id"`
} `json:"message"`
}
if err := json.NewDecoder(response.Body).Decode(&payload); err != nil || payload.Message.ID == "" {
r.sendErrors.Add(1)
return
}
r.mu.Lock()
r.sendTimes[payload.Message.ID] = time.Now()
r.mu.Unlock()
r.sent.Add(1)
}
// clientLoop держит WebSocket, отвечает на HELLO, шлёт heartbeat и меряет
// задержку доставки MESSAGE_CREATE, а также периодически дёргает API.
func (r *runner) clientLoop(ctx context.Context, index int) {
item := r.clients[index]
wsURL := strings.Replace(strings.Replace(r.baseURL, "https://", "wss://", 1), "http://", "ws://", 1) + "/gateway"
conn, handshake, err := websocket.Dial(ctx, wsURL, &websocket.DialOptions{
HTTPHeader: http.Header{"Authorization": []string{"Bearer " + item.Token}},
})
if handshake != nil && handshake.Body != nil {
_ = handshake.Body.Close()
}
if err != nil {
r.disconnect.Add(1)
return
}
defer func() { _ = conn.CloseNow() }()
conn.SetReadLimit(1 << 20)
heartbeat := time.NewTicker(30 * time.Second)
defer heartbeat.Stop()
go func() {
for {
select {
case <-ctx.Done():
return
case <-heartbeat.C:
payload, _ := json.Marshal(envelope{Op: opHeartbeat})
if err := conn.Write(ctx, websocket.MessageText, payload); err != nil {
return
}
}
}
}()
if r.apiEvery > 0 {
go r.apiLoop(ctx, item)
}
for {
kind, data, err := conn.Read(ctx)
if err != nil {
if ctx.Err() == nil {
r.disconnect.Add(1)
}
return
}
if kind != websocket.MessageText {
continue
}
var frame envelope
if err := json.Unmarshal(data, &frame); err != nil {
continue
}
switch frame.Op {
case opHello:
// Токен передаём в IDENTIFY: cookie рукопожатия тоже подошла бы,
// но Bearer-токен — основной путь для внешних клиентов (AGENT.md 8.1).
identify, _ := json.Marshal(map[string]any{"token": item.Token})
payload, _ := json.Marshal(envelope{Op: opIdentify, D: identify})
if err := conn.Write(ctx, websocket.MessageText, payload); err != nil {
return
}
case opHeartbeatAck:
// Ответ на heartbeat: разрыв соединения не нужен.
case opDispatch:
r.received.Add(1)
switch frame.T {
case "READY", "RESUMED":
r.ready.Add(1)
case "MESSAGE_CREATE":
r.noteMessage(frame.D)
}
}
}
}
// noteMessage фиксирует задержку доставки конкретного сообщения.
func (r *runner) noteMessage(payload json.RawMessage) {
var message struct {
ID string `json:"id"`
}
if err := json.Unmarshal(payload, &message); err != nil || message.ID == "" {
return
}
now := time.Now()
r.mu.Lock()
sentAt, ok := r.sendTimes[message.ID]
if ok {
delete(r.sendTimes, message.ID)
r.latencies = append(r.latencies, float64(now.Sub(sentAt).Microseconds())/1000)
}
r.mu.Unlock()
if ok {
r.delivered.Add(1)
}
}
// apiLoop измеряет задержку обычных запросов API тем же клиентом.
func (r *runner) apiLoop(ctx context.Context, item credentials) {
ticker := time.NewTicker(r.apiEvery)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
started := time.Now()
response, err := r.do(ctx, item.Token, http.MethodGet, "/api/v1/users/@me", nil)
if err != nil {
r.apiErrors.Add(1)
continue
}
if response.StatusCode != http.StatusOK {
r.apiErrors.Add(1)
_ = response.Body.Close()
continue
}
_ = response.Body.Close()
r.mu.Lock()
r.apiLatency = append(r.apiLatency, float64(time.Since(started).Microseconds())/1000)
r.mu.Unlock()
r.apiCount.Add(1)
}
}
}
func (r *runner) do(ctx context.Context, token, method, path string, body any) (*http.Response, error) {
var reader *strings.Reader
if body == nil {
reader = strings.NewReader("")
} else {
encoded, err := json.Marshal(body)
if err != nil {
return nil, err
}
reader = strings.NewReader(string(encoded))
}
request, err := http.NewRequestWithContext(ctx, method, r.baseURL+path, reader)
if err != nil {
return nil, err
}
request.Header.Set("Authorization", "Bearer "+token)
if body != nil {
request.Header.Set("Content-Type", "application/json")
}
return r.httpClient().Do(request)
}
func (r *runner) report(started time.Time, duration time.Duration) report {
r.mu.Lock()
latencies := append([]float64(nil), r.latencies...)
apiLatency := append([]float64(nil), r.apiLatency...)
undelivered := len(r.sendTimes)
r.mu.Unlock()
sort.Float64s(latencies)
sort.Float64s(apiLatency)
final := report{
Clients: len(r.clients),
Rate: r.rate,
DurationSeconds: duration.Seconds(),
Sent: r.sent.Load(),
SendErrors: r.sendErrors.Load(),
Received: r.received.Load(),
Delivered: r.delivered.Load(),
APIRequests: r.apiCount.Load(),
APIErrors: r.apiErrors.Load(),
ReadyClients: int(r.ready.Load()),
Disconnects: r.disconnect.Load(),
StartedAt: started,
}
final.DeliveryP50MS = percentile(latencies, 0.50)
final.DeliveryP95MS = percentile(latencies, 0.95)
final.DeliveryP99MS = percentile(latencies, 0.99)
if len(latencies) > 0 {
final.DeliveryMaxMS = latencies[len(latencies)-1]
}
final.APIP50MS = percentile(apiLatency, 0.50)
final.APIP95MS = percentile(apiLatency, 0.95)
if len(apiLatency) > 0 {
final.APIMaxMS = apiLatency[len(apiLatency)-1]
}
if undelivered > 0 {
fmt.Printf("внимание: %d сообщений не дошли до клиентов за время прогона\n", undelivered)
}
return final
}
func percentile(sorted []float64, quantile float64) float64 {
if len(sorted) == 0 {
return 0
}
index := int(quantile * float64(len(sorted)-1))
return sorted[index]
}
func printReport(final report) {
fmt.Printf("\n=== нагрузка: %d клиентов, %.0f msg/s, %.0f с ===\n",
final.Clients, final.Rate, final.DurationSeconds)
fmt.Printf("клиентов с READY: %d\n", final.ReadyClients)
fmt.Printf("отправлено сообщений: %d (ошибок %d)\n", final.Sent, final.SendErrors)
fmt.Printf("доставлено: %d\n", final.Delivered)
fmt.Printf("доставка, мс: p50 %.1f | p95 %.1f | p99 %.1f | max %.1f\n",
final.DeliveryP50MS, final.DeliveryP95MS, final.DeliveryP99MS, final.DeliveryMaxMS)
fmt.Printf("API /users/@me, мс: p50 %.1f | p95 %.1f | max %.1f (запросов %d, ошибок %d)\n",
final.APIP50MS, final.APIP95MS, final.APIMaxMS, final.APIRequests, final.APIErrors)
fmt.Printf("событий получено: %d, разрывов WS: %d\n", final.Received, final.Disconnects)
}
func writeReport(path string, final report) error {
encoded, err := json.MarshalIndent(final, "", " ")
if err != nil {
return err
}
// Отчёт не содержит секретов, но рядом лежат учётные данные: 600.
return os.WriteFile(path, encoded, 0o600)
}