// 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) }