From 66dd07eb7b6696f4dc6858869a47b8a611fa9738 Mon Sep 17 00:00:00 2001 From: grendervill Date: Tue, 22 Sep 2026 20:00:33 +0300 Subject: [PATCH] =?UTF-8?q?feat(ops):=20=D0=B3=D0=B5=D0=BD=D0=B5=D1=80?= =?UTF-8?q?=D0=B0=D1=82=D0=BE=D1=80=20=D0=BD=D0=B0=D0=B3=D1=80=D1=83=D0=B7?= =?UTF-8?q?=D0=BA=D0=B8=20=D0=B8=20=D0=BF=D0=BE=D0=B4=D0=B3=D0=BE=D1=82?= =?UTF-8?q?=D0=BE=D0=B2=D0=BA=D0=B0=20=D0=B0=D0=BA=D0=BA=D0=B0=D1=83=D0=BD?= =?UTF-8?q?=D1=82=D0=BE=D0=B2=20=D0=B4=D0=BB=D1=8F=20=D1=82=D0=B5=D1=81?= =?UTF-8?q?=D1=82=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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*. --- .gitignore | 2 +- cmd/glchat/cli.go | 105 +++++++ cmd/glchat/main.go | 2 +- cmd/loadgen/main.go | 577 ++++++++++++++++++++++++++++++++++ deploy/glchat-cli.sh | 10 + internal/auth/service.go | 55 ++++ internal/auth/service_test.go | 45 +++ web/e2e/voice-ten.spec.ts | 232 ++++++++++++++ 8 files changed, 1026 insertions(+), 2 deletions(-) create mode 100644 cmd/loadgen/main.go create mode 100644 web/e2e/voice-ten.spec.ts diff --git a/.gitignore b/.gitignore index 7a18207..cb10f56 100644 --- a/.gitignore +++ b/.gitignore @@ -38,5 +38,5 @@ web/node_modules/.cache/ !.vscode/extensions.json # артефакты Playwright (e2e-прогоны) — локальные -web/test-results/ +web/test-results*/ web/playwright-report/ diff --git a/cmd/glchat/cli.go b/cmd/glchat/cli.go index ea4b85d..a34bb8b 100644 --- a/cmd/glchat/cli.go +++ b/cmd/glchat/cli.go @@ -2,11 +2,13 @@ package main import ( "context" + "encoding/json" "errors" "flag" "fmt" "log/slog" "os" + "strconv" "time" "github.com/pquerna/otp/totp" @@ -33,6 +35,8 @@ func runCLI(args []string) int { return cliTOTPSetup(args[1:]) case "totp-reset": return cliTOTPReset(args[1:]) + case "create-user": + return cliCreateUser(args[1:]) case "cleanup": return cliCleanup(args[1:]) case "reindex": @@ -273,6 +277,107 @@ func cliSetInstanceAdmin(command string, args []string) int { // cliCleanup запускает чистку вручную: сессии, аудит и файлы без ссылок // (AGENT.md 6.4, 10.7). +// cliCreateUser создаёт аккаунты по команде обслуживания: регистрация на +// инстансе может быть выключена, а для нагрузочного теста (AGENT.md 11.4) +// нужны десятки пользователей. Пароль задаёт оператор, аккаунты создаются без +// сессий; список логинов и паролей можно выгрузить в файл для генератора +// нагрузки (--out). +func cliCreateUser(args []string) int { + flags := flag.NewFlagSet("create-user", flag.ContinueOnError) + email := flags.String("email", "", "email пользователя") + username := flags.String("username", "", "username (по умолчанию — часть email)") + password := flags.String("password", "", "пароль") + count := flags.Int("count", 1, "сколько аккаунтов создать (для нагрузки)") + prefix := flags.String("prefix", "", "префикс логина в пакетном режиме, например loadtest") + domain := flags.String("domain", "loadtest.local", "домен email в пакетном режиме") + out := flags.String("out", "", "файл с созданными учётными данными (JSON)") + sessions := flags.Bool("sessions", false, "выдать каждому аккаунту сессию (для нагрузочного теста)") + if err := flags.Parse(args); err != nil { + return 2 + } + if *password == "" { + fmt.Fprintln(os.Stderr, "glchat create-user: укажите --password") + return 2 + } + if *count > 1 && *prefix == "" { + fmt.Fprintln(os.Stderr, "glchat create-user: для --count > 1 нужен --prefix") + return 2 + } + if *count <= 1 && *email == "" { + fmt.Fprintln(os.Stderr, "glchat create-user: укажите --email") + return 2 + } + + type credentials struct { + ID string `json:"id"` + Username string `json:"username"` + Email string `json:"email"` + Password string `json:"password"` + Token string `json:"token,omitempty"` + } + created := make([]credentials, 0, *count) + + code := withServices(func(ctx context.Context, _ *store.Store, authService *auth.Service, _ config.Config) error { + for index := 1; index <= *count; index += 1 { + name := *username + address := *email + if *count > 1 { + name = fmt.Sprintf("%s%d", *prefix, index) + address = fmt.Sprintf("%s%d@%s", *prefix, index, *domain) + } + user, err := authService.CreateUserByOperator(ctx, auth.RegisterInput{ + Username: name, + Email: address, + Password: *password, + }) + if err != nil { + if errors.Is(err, auth.ErrUsernameTaken) || errors.Is(err, auth.ErrEmailTaken) { + fmt.Printf("пропущен (уже существует): %s\n", address) + continue + } + return fmt.Errorf("создать %s: %w", address, err) + } + entry := credentials{ + ID: strconv.FormatUint(user.ID, 10), + Username: user.Username, + Email: address, + Password: *password, + } + if *sessions { + token, err := authService.IssueSessionForOperator(ctx, user.ID, "glchat-cli loadtest") + if err != nil { + return fmt.Errorf("сессия для %s: %w", address, err) + } + entry.Token = token + } + created = append(created, entry) + if *count == 1 || index%25 == 0 || index == *count { + fmt.Printf("создано %d из %d\n", index, *count) + } + } + return nil + }) + if code != 0 { + return code + } + if *out != "" { + // #nosec G117 -- файл с тестовыми паролями создаётся намеренно, права 600. + encoded, err := json.MarshalIndent(created, "", " ") + if err != nil { + fmt.Fprintf(os.Stderr, "glchat create-user: %v\n", err) + return 1 + } + // Файл содержит пароли тестовых аккаунтов: права 600. + if err := os.WriteFile(*out, encoded, 0o600); err != nil { + fmt.Fprintf(os.Stderr, "glchat create-user: %v\n", err) + return 1 + } + fmt.Printf("учётные данные: %s (%d)\n", *out, len(created)) + } + fmt.Printf("создано аккаунтов: %d\n", len(created)) + return 0 +} + func cliCleanup(args []string) int { flags := flag.NewFlagSet("cleanup", flag.ContinueOnError) auditDays := flags.Int("audit-days", -1, "сколько дней хранить журнал аудита (по умолчанию из конфигурации)") diff --git a/cmd/glchat/main.go b/cmd/glchat/main.go index 07204b2..2614a39 100644 --- a/cmd/glchat/main.go +++ b/cmd/glchat/main.go @@ -43,7 +43,7 @@ func main() { if len(os.Args) > 1 { switch os.Args[1] { case "bootstrap-admin", "reset-password", "make-admin", "remove-admin", "totp-setup", "totp-reset", - "cleanup", "reindex": + "create-user", "cleanup", "reindex": os.Exit(runCLI(os.Args[1:])) } } diff --git a/cmd/loadgen/main.go b/cmd/loadgen/main.go new file mode 100644 index 0000000..97db48b --- /dev/null +++ b/cmd/loadgen/main.go @@ -0,0 +1,577 @@ +// 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) +} diff --git a/deploy/glchat-cli.sh b/deploy/glchat-cli.sh index ef1cd13..3f620d4 100755 --- a/deploy/glchat-cli.sh +++ b/deploy/glchat-cli.sh @@ -29,6 +29,7 @@ glchat — управление инстансом glchat security-check аудит: SSH, fail2ban, порты, лимиты glchat reindex пересобрать индекс поиска (FTS5) glchat cleanup чистка: сессии, аудит, файлы без ссылок + glchat create-user создать аккаунты (--count/--prefix/--sessions) glchat change-domain сменить домен glchat uninstall [--keep-data] удалить стек glchat version версия CLI и приложения @@ -175,6 +176,14 @@ cmd_cleanup() { ${COMPOSE} exec -T app glchat cleanup "$@" } +# cmd_create_user создаёт аккаунты командой обслуживания: регистрация может +# быть выключена, а для нагрузочного теста нужны десятки пользователей +# (AGENT.md 10.7, 11.4). Флаг --sessions сразу выдаёт сессии для генератора. +cmd_create_user() { + require_stack + ${COMPOSE} exec -T app glchat create-user "$@" +} + main() { local cmd="${1:-help}" shift || true @@ -196,6 +205,7 @@ main() { reindex) cmd_reindex "$@" ;; verify-backup) cmd_verify_backup "$@" ;; cleanup) cmd_cleanup "$@" ;; + create-user) cmd_create_user "$@" ;; change-domain) cmd_change_domain "$@" ;; uninstall) cmd_uninstall "$@" ;; version) cmd_version "$@" ;; diff --git a/internal/auth/service.go b/internal/auth/service.go index 97a1d3b..a25cd33 100644 --- a/internal/auth/service.go +++ b/internal/auth/service.go @@ -213,6 +213,61 @@ func (s *Service) Register(ctx context.Context, in RegisterInput) (*store.User, return user, token, session, nil } +// CreateUserByOperator создаёт аккаунт по команде обслуживания (CLI +// create-user): регистрация на инстансе может быть выключена, а сессия +// создателю не нужна — он действует от имени администратора. Используется в +// том числе для подготовки нагрузки (AGENT.md 11.4). +func (s *Service) CreateUserByOperator(ctx context.Context, in RegisterInput) (*store.User, error) { + username := strings.TrimSpace(in.Username) + if err := crypto.ValidateUsername(username); err != nil { + return nil, ErrInvalidUsername + } + if err := s.policy.validate(in.Password); err != nil { + return nil, ErrWeakPassword + } + email := crypto.NormalizeEmail(in.Email) + if err := validateEmail(email); err != nil { + return nil, err + } + passwordHash, err := s.hasher.Hash(in.Password) + if err != nil { + return nil, fmt.Errorf("hash password: %w", err) + } + userID := s.store.NextID() + emailEncrypted, err := s.encryptEmail(userID, email) + if err != nil { + return nil, err + } + user, err := s.store.CreateUser(ctx, store.CreateUserParams{ + ID: userID, + Username: username, + DisplayName: defaultDisplayName(in.DisplayName, username), + EmailEnc: emailEncrypted, + EmailIndex: s.masterKey.BlindIndex(email), + PasswordHash: passwordHash, + Locale: in.Locale, + }) + if err != nil { + if errors.Is(err, store.ErrConflict) { + if existing, lookupErr := s.store.GetUserByUsername(ctx, username); lookupErr == nil && existing != nil { + return nil, ErrUsernameTaken + } + return nil, ErrEmailTaken + } + return nil, fmt.Errorf("create user: %w", err) + } + _ = s.store.RecordSecurityEvent(ctx, &user.ID, "created_by_operator", "", "", "") + return user, nil +} + +// IssueSessionForOperator выдаёт сессию по команде обслуживания: нужна для +// нагрузочного теста (AGENT.md 11.4), где сотни клиентов не могут пройти вход +// через /auth/login — лимит попыток по IP рассчитан на защиту от перебора. +func (s *Service) IssueSessionForOperator(ctx context.Context, userID uint64, userAgent string) (string, error) { + token, _, err := s.createSession(ctx, userID, userAgent, "") + return token, err +} + // UserByEmail ищет аккаунт по email через blind index (без расшифровки БД). func (s *Service) UserByEmail(ctx context.Context, email string) (*store.User, error) { return s.userByEmail(ctx, email) diff --git a/internal/auth/service_test.go b/internal/auth/service_test.go index a7e8e84..6fb294c 100644 --- a/internal/auth/service_test.go +++ b/internal/auth/service_test.go @@ -459,3 +459,48 @@ func totpCode(t *testing.T, secret string) string { } return code } + +// TestCreateUserByOperator проверяет команду обслуживания: аккаунт создаётся +// даже при выключенной регистрации, дубликаты и слабые пароли отклоняются +// (AGENT.md 10.7, 11.4 — подготовка нагрузки). +func TestCreateUserByOperator(t *testing.T) { + ctx := context.Background() + service, st := newService(t) + + if err := st.SetInstanceSetting(ctx, "registration_enabled", "false"); err != nil { + t.Fatalf("SetInstanceSetting: %v", err) + } + // Публичная регистрация закрыта. + if _, _, _, err := service.Register(ctx, auth.RegisterInput{ + Username: "operator_blocked", Email: "blocked@example.com", Password: "correct-horse-battery", + }); !errors.Is(err, auth.ErrRegistrationOff) { + t.Fatalf("Register with registration off = %v, want ErrRegistrationOff", err) + } + + user, err := service.CreateUserByOperator(ctx, auth.RegisterInput{ + Username: "loadtest1", Email: "loadtest1@loadtest.local", Password: "correct-horse-battery", + }) + if err != nil { + t.Fatalf("CreateUserByOperator: %v", err) + } + if user.Username != "loadtest1" || user.IsInstanceAdmin { + t.Fatalf("created user = %+v", user) + } + // Вход новым аккаунтом работает без выдачи сессии создателю. + if _, _, _, err := service.Login(ctx, auth.LoginInput{ + Email: "loadtest1@loadtest.local", Password: "correct-horse-battery", + }); err != nil { + t.Fatalf("Login as created user: %v", err) + } + + if _, err := service.CreateUserByOperator(ctx, auth.RegisterInput{ + Username: "loadtest1", Email: "other@loadtest.local", Password: "correct-horse-battery", + }); !errors.Is(err, auth.ErrUsernameTaken) { + t.Fatalf("duplicate username = %v, want ErrUsernameTaken", err) + } + if _, err := service.CreateUserByOperator(ctx, auth.RegisterInput{ + Username: "loadtest2", Email: "loadtest2@loadtest.local", Password: "short", + }); !errors.Is(err, auth.ErrWeakPassword) { + t.Fatalf("weak password = %v, want ErrWeakPassword", err) + } +} diff --git a/web/e2e/voice-ten.spec.ts b/web/e2e/voice-ten.spec.ts new file mode 100644 index 0000000..66ee069 --- /dev/null +++ b/web/e2e/voice-ten.spec.ts @@ -0,0 +1,232 @@ +import { readFileSync } from 'node:fs'; +import { expect, test, type Browser, type BrowserContext, type Page } from '@playwright/test'; + +/** + * Голосовая нагрузка (AGENT.md 9.1, 11.4): десять участников в одной комнате + * и внешний 1080p60-поток. Сессии берутся из файла, который готовит команда + * обслуживания (`glchat create-user --sessions`), — вход через форму упёрся бы + * в лимит попыток по IP, рассчитанный на защиту от перебора. + * + * GLCHAT_E2E_VOICE_TEN=1 GLCHAT_URL=https://gl.mhspx.su \ + * GLCHAT_LOADTEST_CREDS=../../build/loadtest-users.json \ + * GLCHAT_LOADTEST_FIXTURE=../../build/loadtest-fixture.json \ + * npx playwright test e2e/voice-ten.spec.ts + */ +const enabled = process.env.GLCHAT_E2E_VOICE_TEN === '1'; +const CREDS_PATH = process.env.GLCHAT_LOADTEST_CREDS ?? '../../build/loadtest-users.json'; +const FIXTURE_PATH = process.env.GLCHAT_LOADTEST_FIXTURE ?? '../../build/loadtest-fixture.json'; +const PARTICIPANTS = Number(process.env.GLCHAT_VOICE_PARTICIPANTS ?? '10'); +const EXPECT_SCREENSHARE = process.env.GLCHAT_EXPECT_SCREENSHARE === '1'; +// Идентификатор внешнего издателя: реальный участник сервера (см. build/loadtest-voice.sh). +const SCREENSHARE_IDENTITY = process.env.GLCHAT_SCREENSHARE_IDENTITY ?? 'screenshare-1080p60'; + +interface Account { + id: string; + username: string; + token: string; +} + +interface Fixture { + guildId: string; + voiceChannelId: string; +} + +function readAccounts(): Account[] { + const raw = JSON.parse(readFileSync(CREDS_PATH, 'utf8')) as Account[]; + return raw.filter((item) => item.token !== ''); +} + +function readFixture(): Fixture { + return JSON.parse(readFileSync(FIXTURE_PATH, 'utf8')) as Fixture; +} + +/** signInWithToken открывает контекст с готовой сессией: форма входа не нужна. */ +async function signInWithToken( + browser: Browser, + account: Account, + baseURL: string, +): Promise<{ context: BrowserContext; page: Page }> { + const context = await browser.newContext({ + permissions: ['microphone', 'camera'], + baseURL, + }); + await context.addCookies([ + // Cookie сессии — __Host-session (internal/gateway/session.go). + { name: '__Host-session', value: account.token, url: baseURL, httpOnly: true, secure: true }, + ]); + // Аккаунты создаёт команда обслуживания, онбординг у них не пройден: без + // этого AuthGuard уводит браузер на /onboarding и комнаты не видно. + const headers = { Authorization: `Bearer ${account.token}` }; + const me = await context.request.get('/api/v1/users/@me', { headers }); + if (me.ok()) { + const payload = (await me.json()) as { + user: { onboarding_completed: boolean; display_name?: string | null; username: string }; + }; + if (!payload.user.onboarding_completed) { + await context.request.post('/api/v1/users/@me/onboarding/complete', { + headers, + data: { display_name: payload.user.display_name ?? payload.user.username }, + }); + } + } + const page = await context.newPage(); + await page.goto('/app'); + await expect(page).toHaveURL(/\/app/u, { timeout: 30_000 }); + return { context, page }; +} + +/** joinVoice нажимает «Войти» в голосовой комнате и ждёт подключения. */ +async function joinVoice(page: Page): Promise { + const button = page.getByTestId('voice-join'); + await button.waitFor({ state: 'visible', timeout: 30_000 }); + await button.click({ force: true, timeout: 10_000 }).catch(() => undefined); + await expect(button).toBeHidden({ timeout: 60_000 }); +} + +/** inboundAudioBytes суммирует полученные аудио-байты по всем соединениям. */ +async function inboundAudioBytes(page: Page): Promise { + return page.evaluate(async () => { + const globalWindow = window as unknown as { __pcs?: RTCPeerConnection[] }; + let total = 0; + for (const pc of globalWindow.__pcs ?? []) { + const stats = await pc.getStats(); + stats.forEach((raw) => { + const report = raw as Record; + if (report['type'] === 'inbound-rtp' && report['kind'] === 'audio') { + total += Number(report['bytesReceived'] ?? 0); + } + }); + } + return total; + }); +} + +/** inboundVideoStats отдаёт лучший входящий видеопоток: размер и кадры в секунду. */ +async function inboundVideoStats( + page: Page, +): Promise<{ width: number; height: number; fps: number }> { + return page.evaluate(async () => { + const globalWindow = window as unknown as { __pcs?: RTCPeerConnection[] }; + let best = { width: 0, height: 0, fps: 0, decoded: 0, bytes: 0 }; + for (const pc of globalWindow.__pcs ?? []) { + const stats = await pc.getStats(); + stats.forEach((raw) => { + const report = raw as Record; + if (report['type'] !== 'inbound-rtp' || report['kind'] !== 'video') { + return; + } + const width = Number(report['frameWidth'] ?? 0); + const height = Number(report['frameHeight'] ?? 0); + const fps = Number(report['framesPerSecond'] ?? 0); + const decoded = Number(report['framesDecoded'] ?? 0); + const bytes = Number(report['bytesReceived'] ?? 0); + if (width * height > best.width * best.height) { + best = { width, height, fps, decoded, bytes }; + } + }); + } + return best; + }); +} + +/** instrumentPeerConnections собирает соединения страницы для статистики. */ +async function instrumentPeerConnections(context: BrowserContext): Promise { + await context.addInitScript(() => { + const globalWindow = window as unknown as { __pcs?: RTCPeerConnection[] }; + globalWindow.__pcs = []; + const Original = window.RTCPeerConnection; + window.RTCPeerConnection = function Patched( + ...args: ConstructorParameters + ) { + const pc = new Original(...args); + globalWindow.__pcs?.push(pc); + return pc; + } as unknown as typeof RTCPeerConnection; + window.RTCPeerConnection.prototype = Original.prototype; + }); +} + +test.describe('голосовая нагрузка: десять участников', () => { + test.skip(!enabled, 'нужен развёрнутый инстанс: GLCHAT_E2E_VOICE_TEN=1'); + test('десять клиентов видят друг друга, аудио и 1080p60 доходят', async ({ browser }) => { + test.setTimeout(300_000); + const baseURL = process.env.GLCHAT_URL ?? 'https://gl.mhspx.su'; + const accounts = readAccounts().slice(0, PARTICIPANTS); + const fixture = readFixture(); + expect(accounts.length, 'нужны аккаунты с сессиями').toBeGreaterThanOrEqual(PARTICIPANTS); + + const sessions: { account: Account; context: BrowserContext; page: Page }[] = []; + try { + for (const account of accounts) { + const context = await browser.newContext({ + permissions: ['microphone', 'camera'], + baseURL, + }); + await instrumentPeerConnections(context); + await context.addCookies([ + { + name: '__Host-session', + value: account.token, + url: baseURL, + httpOnly: true, + secure: true, + }, + ]); + const page = await context.newPage(); + await page.goto(`/app/${fixture.guildId}/${fixture.voiceChannelId}`); + await expect(page).toHaveURL(/\/app/u, { timeout: 30_000 }); + sessions.push({ account, context, page }); + } + + // Все десять входят в комнату. + for (const session of sessions) { + await joinVoice(session.page); + } + console.log(`участников в комнате: ${sessions.length}`); + + // Каждый видит плитки всех участников: это состояние LiveKit, а не список сервера. + for (const session of sessions.slice(0, 2)) { + for (const other of sessions) { + await expect(session.page.getByTestId(`voice-tile-${other.account.id}`)).toBeVisible({ + timeout: 60_000, + }); + } + } + + // Аудио реально доходит: у двух клиентов считаем входящие аудио-байты. + for (const session of sessions.slice(0, 2)) { + await expect + .poll(async () => inboundAudioBytes(session.page), { timeout: 90_000, intervals: [1000] }) + .toBeGreaterThan(0); + } + + if (EXPECT_SCREENSHARE) { + // Внешний 1080p60-поток (lk room join с IVF) должен дойти до клиента. + // Adaptive stream ставит видео на паузу, пока плитка не на экране, а + // плитки внешнего издателя в интерфейсе может не быть вовсе (он не + // проходит через приложение). Поэтому здесь проверяется, что слой + // 1080p60 доехал до клиента (размер кадра и байты), а реальное + // декодирование 60 кадров в секунду меряет независимый подписчик — + // build/verify-1080p60-browser.mjs (adaptive stream выключен). + await expect + .poll(async () => inboundVideoStats(sessions[0]!.page), { + timeout: 90_000, + intervals: [2000], + }) + .toMatchObject({ width: 1920, height: 1080 }); + const stats = await inboundVideoStats(sessions[0]!.page); + console.log('входящее видео:', JSON.stringify(stats)); + await expect + .poll(async () => (await inboundVideoStats(sessions[0]!.page)).bytes, { + timeout: 60_000, + intervals: [2000], + }) + .toBeGreaterThan(0); + } + } finally { + for (const session of sessions) { + await session.context.close(); + } + } + }); +});