578 lines
18 KiB
Go
578 lines
18 KiB
Go
|
|
// 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)
|
|||
|
|
}
|