/
t3
/
cli
Обзор
Документация
Войти
/
t3
/
cli
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
internal/pool/pool.go
312 строк
8 KB
Ivan
move s3 to the root
09 июл 2026, 13:51
09 июл 2026, 13:51
c214eff
Код
Авторство
О чём код?
package pool import ( "context" "crypto/tls" "encoding/json" "fmt" "log/slog" "net/http" "sync" "sync/atomic" "time" "github.com/dop251/goja" "gitverse.ru/t3/cli/internal/engine" "gitverse.ru/t3/cli/internal/engine/bridge" "gitverse.ru/t3/cli/internal/metrics" ) type Pool struct { ctx context.Context cancel context.CancelFunc lastUserID atomic.Int64 logger *slog.Logger en *engine.Engine setupData any s3Cache *bridge.S3Cache s3Config *bridge.S3Config agentData goja.Value // результат agent_init() activeUsers []*User stoppedUsers []*User mx sync.Mutex wg *sync.WaitGroup exitErr error registry *metrics.MetricsRegistry thinkTimeMs int64 pacingMs int64 httpTransport *http.Transport } func NewPool(ctx context.Context, logger *slog.Logger, bundle, setupData []byte, registry *metrics.MetricsRegistry, thinkTimeMs, pacingMs int64, cacheDir string, cacheMaxSize int64) (*Pool, error) { en, err := engine.NewEngine(bundle) if err != nil { return nil, err } var setupDataObj any if len(setupData) > 0 { err = json.Unmarshal(setupData, &setupDataObj) if err != nil { return nil, err } } // Извлекаем S3 конфиг из setup_data s3Config := extractS3Config(setupDataObj) // Создаем S3 кэш (с поддержкой дискового кэша между запусками) s3Cache := bridge.NewS3Cache(5*time.Minute, cacheDir, cacheMaxSize) poolCtx, cancel := context.WithCancel(ctx) pool := &Pool{ ctx: poolCtx, cancel: cancel, en: en, logger: logger, setupData: setupDataObj, s3Cache: s3Cache, s3Config: s3Config, activeUsers: make([]*User, 0, 500), stoppedUsers: make([]*User, 0, 500), wg: &sync.WaitGroup{}, registry: registry, thinkTimeMs: thinkTimeMs, pacingMs: pacingMs, httpTransport: newHTTPTransport(), } // Запускаем agent_init() — один раз на весь пул, до создания VUser if err := pool.runInitAgent(); err != nil { return nil, fmt.Errorf("agent_init failed: %w", err) } return pool, nil } // extractS3Config извлекает S3 конфигурацию из setup_data. // Ожидается структура: { "s3": { "endpoint": "...", "accessKey": "...", ... } } // Также поддерживает ключ "s3Config" (для совместимости с разными naming conventions). func extractS3Config(setupData any) *bridge.S3Config { data, ok := setupData.(map[string]interface{}) if !ok { return nil } s3Raw, ok := data["s3"] if !ok { // Fallback: пробуем "s3Config" как альтернативное имя ключа s3Raw, ok = data["s3Config"] if !ok { return nil } } s3Map, ok := s3Raw.(map[string]interface{}) if !ok { return nil } cfg := &bridge.S3Config{ Region: "ru-west-1", } if v, ok := s3Map["endpoint"]; ok { cfg.Endpoint = fmt.Sprint(v) } if v, ok := s3Map["accessKey"]; ok { cfg.AccessKey = fmt.Sprint(v) } if v, ok := s3Map["secretKey"]; ok { cfg.SecretKey = fmt.Sprint(v) } if v, ok := s3Map["region"]; ok { cfg.Region = fmt.Sprint(v) } if v, ok := s3Map["useSSL"]; ok { if b, ok := v.(bool); ok { cfg.UseSSL = b } } if v, ok := s3Map["insecureSkipVerify"]; ok { if b, ok := v.(bool); ok { cfg.InsecureSkipVerify = b } } if v, ok := s3Map["basePath"]; ok { cfg.BasePath = fmt.Sprint(v) } return cfg } // runInitAgent создает временный Instance и вызывает agent_init(setupData). // После возврата запечатывает S3 кэш. func (pool *Pool) runInitAgent() error { // Создаем временный Instance для выполнения agent_init initCfg := &engine.InstanceConfig{ Logger: pool.logger, Metrics: pool.registry, S3Cache: pool.s3Cache, S3Cfg: pool.s3Config, Ctx: pool.ctx, HTTPTransport: pool.httpTransport, } instance, err := pool.en.NewInstance(initCfg) if err != nil { return fmt.Errorf("cannot create agent_init instance: %w", err) } initAgentFunc, ok := instance.GetFunc("agent_init") if ok { setupDataValue := instance.GetRuntime().ToValue(pool.setupData) result, callErr := initAgentFunc.Call(setupDataValue) if callErr != nil { return fmt.Errorf("agent_init execution failed: %w", callErr) } // Сохраняем результат agent_init для передачи в VUser pool.agentData = result } // Запечатываем кэш — теперь download() будет только читать из кэша pool.s3Cache.Seal() return nil } // newHTTPTransport создаёт shared http.Transport, оптимизированный для нагрузочного тестирования. // Соединения переиспользуются, лимиты на хост высокие, keepalive включён. func newHTTPTransport() *http.Transport { return &http.Transport{ MaxIdleConns: 0, // не ограничено MaxIdleConnsPerHost: 500, MaxConnsPerHost: 0, // не ограничено IdleConnTimeout: 90 * time.Second, TLSHandshakeTimeout: 10 * time.Second, ResponseHeaderTimeout: 30 * time.Second, DisableKeepAlives: false, DisableCompression: false, TLSClientConfig: &tls.Config{ InsecureSkipVerify: false, }, } } func (pool *Pool) Stop(err error) { pool.mx.Lock() defer pool.mx.Unlock() if err != nil && pool.exitErr == nil { pool.exitErr = err } pool.cancel() } func (pool *Pool) Wait() error { pool.wg.Wait() return pool.exitErr } func (pool *Pool) GoTo(target int, duration time.Duration) error { if !pool.mx.TryLock() { return fmt.Errorf("busy") } defer pool.mx.Unlock() if target < 0 { return fmt.Errorf("target cannot be a negative number") } activeUsers := len(pool.activeUsers) if target < activeUsers { diff := activeUsers - target return pool.rampDown(diff, duration) } if target > activeUsers { diff := target - activeUsers return pool.rampUp(diff, duration) } return nil } func (pool *Pool) rampUp(users int, duration time.Duration) error { nanoseconds := duration.Nanoseconds() / int64(users) for range users { select { case <-pool.ctx.Done(): return nil case <-time.After(time.Duration(nanoseconds)): user, err := pool.allocateUser() if err != nil { return err } if pool.registry != nil { pool.registry.IncActiveUsers(1) } pool.wg.Go(func() { if err := user.Run(pool.ctx); err != nil { pool.Stop(err) } }) } } return nil } func (pool *Pool) rampDown(users int, duration time.Duration) error { nanoseconds := duration.Nanoseconds() / int64(users) for range users { select { case <-pool.ctx.Done(): return nil case <-time.After(time.Duration(nanoseconds)): pool.stopUser() } } return nil } func (pool *Pool) stopUser() { if len(pool.activeUsers) > 0 { user := pool.activeUsers[0] user.instance.Interrupt("ramp down") pool.activeUsers = pool.activeUsers[1:] pool.stoppedUsers = append(pool.stoppedUsers, user) if pool.registry != nil { pool.registry.DecActiveUsers(1) } } } func (pool *Pool) allocateUser() (*User, error) { if len(pool.stoppedUsers) > 0 { user := pool.stoppedUsers[0] pool.stoppedUsers = pool.stoppedUsers[1:] user.instance.ClearInterrupt() pool.activeUsers = append(pool.activeUsers, user) return user, nil } user, err := NewUser( pool.lastUserID.Add(1), pool.logger, pool.en, pool.setupData, pool.registry, pool.s3Cache, pool.s3Config, pool.agentData, pool.thinkTimeMs, pool.pacingMs, pool.httpTransport, ) if err != nil { return nil, fmt.Errorf("cannot create a new user: %w", err) } pool.activeUsers = append(pool.activeUsers, user) return user, nil }