/
t3
/
cli
Обзор
Документация
Войти
/
t3
/
cli
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
internal/pool/user.go
223 строки
6 KB
Ivan
fix user_init args
09 июл 2026, 14:05
09 июл 2026, 14:05
c30c424
Код
Авторство
О чём код?
package pool import ( "context" "fmt" "log/slog" "net/http" "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 User struct { id int64 iteration int64 logger *slog.Logger instance *engine.Instance setupData goja.Value initData goja.Value agentData goja.Value defaultFunc *engine.Func registry *metrics.MetricsRegistry thinkTimeMs int64 pacingMs int64 } func NewUser(id int64, logger *slog.Logger, en *engine.Engine, setupData any, registry *metrics.MetricsRegistry, s3Cache *bridge.S3Cache, s3Cfg *bridge.S3Config, agentData goja.Value, thinkTimeMs, pacingMs int64, httpTransport *http.Transport) (*User, error) { cfg := &engine.InstanceConfig{ Logger: logger, Metrics: registry, S3Cache: s3Cache, S3Cfg: s3Cfg, Ctx: context.Background(), HTTPTransport: httpTransport, } instance, err := en.NewInstance(cfg) if err != nil { return nil, err } setupDataValue := instance.GetRuntime().ToValue(setupData) var initResult goja.Value initFunc, ok := instance.GetFunc("user_init") if ok { // Передаём в user_init один аргумент: // - результат agent_init(), если он определён // - результат setup() в противном случае (agent_init не вызывался) var userInitArg goja.Value if goja.IsUndefined(agentData) || goja.IsNull(agentData) { userInitArg = setupDataValue } else { userInitArg = agentData } initResult, err = initFunc.Call(userInitArg) if err != nil { return nil, err } } defaultFunc, ok := instance.GetFunc("default") if !ok { return nil, fmt.Errorf("'default' function not found") } return &User{ id: id, iteration: 0, logger: logger, instance: instance, defaultFunc: defaultFunc, setupData: setupDataValue, initData: initResult, agentData: agentData, registry: registry, thinkTimeMs: thinkTimeMs, pacingMs: pacingMs, }, nil } func (user *User) Run(ctx context.Context) error { for { select { case <-ctx.Done(): return nil default: } var sleepDuration time.Duration if user.pacingMs > 0 { // Pacing: fixed interval between iteration STARTs iterStart := time.Now() err := user.RunOnce() if user.registry != nil { user.registry.IncTotalIterations() } if err != nil { // InterruptedError is not a real error — happens during ramp-down if _, ok := err.(*goja.InterruptedError); ok { return nil } if user.registry != nil { user.registry.IncErrorIterations() } if _, ok := err.(*goja.Exception); ok { return fmt.Errorf("js exception: %w", err) } if _, ok := err.(*goja.StackOverflowError); ok { return fmt.Errorf("js stackoverflow: %w", err) } user.logger.Error("iteration failed", "user", user.id, "iteration", user.currentIteration(), "error", err.Error()) } elapsed := time.Since(iterStart) sleepDuration = time.Duration(user.pacingMs)*time.Millisecond - elapsed if sleepDuration < 0 { sleepDuration = 0 } } else { // Static think_time: fixed delay after each iteration err := user.RunOnce() if user.registry != nil { user.registry.IncTotalIterations() } if err != nil { // InterruptedError is not a real error — happens during ramp-down if _, ok := err.(*goja.InterruptedError); ok { return nil } if user.registry != nil { user.registry.IncErrorIterations() } if _, ok := err.(*goja.Exception); ok { return fmt.Errorf("js exception: %w", err) } if _, ok := err.(*goja.StackOverflowError); ok { return fmt.Errorf("js stackoverflow: %w", err) } user.logger.Error("iteration failed", "user", user.id, "iteration", user.currentIteration(), "error", err.Error()) } sleepDuration = time.Duration(user.thinkTimeMs) * time.Millisecond } // Wait for either the sleep to complete or context to be cancelled select { case <-ctx.Done(): return nil case <-time.After(sleepDuration): } } } func (user *User) RunOnce() (err error) { defer func() { if r := recover(); r != nil { user.logger.Error("panic in user iteration", "user", user.id, "iteration", user.currentIteration(), "panic", r) err = fmt.Errorf("panic in user iteration: %v", r) } }() // Build the data object by merging init chain with priority (last wins): // 1. setup() result (lowest priority) // 2. agent_init() result // 3. user_init() result (highest priority) // Then overlay system fields that must not be overridden. data := user.instance.GetRuntime().NewObject() mergeInto(data, user.setupData) mergeInto(data, user.agentData) mergeInto(data, user.initData) data.Set("userId", user.id) data.Set("iteration", user.nextIteration()) data.Set("init", user.initData) data.Set("agent", user.agentData) _, err = user.defaultFunc.Call(data) if err != nil { return err } return nil } // mergeInto copies all enumerable properties from src to dst. // If src is nil, undefined, or not a JS object, mergeInto is a no-op. func mergeInto(dst *goja.Object, src goja.Value) { if goja.IsUndefined(src) || goja.IsNull(src) { return } srcObj, ok := src.(*goja.Object) if !ok { // Primitives cannot be iterated; nothing to merge. return } for _, key := range srcObj.Keys() { val := srcObj.Get(key) if !goja.IsUndefined(val) { _ = dst.Set(key, val) } } } func (user *User) currentIteration() int64 { return atomic.LoadInt64(&user.iteration) } func (user *User) nextIteration() int64 { return atomic.AddInt64(&user.iteration, 1) }