From ec2ae8838aac3e4bfc9a7c53520cf2178997dbc4 Mon Sep 17 00:00:00 2001 From: youfak Date: Sun, 2 Aug 2026 08:21:15 +0800 Subject: [PATCH] feat: support load generator request bodies --- README.md | 15 +++++++++++++-- cmd/proxy-loadgen/main.go | 3 ++- cmd/proxy-loadgen/main_test.go | 4 ++-- docs/design/architecture.md | 4 ++-- docs/design/project-structure.md | 2 +- docs/development/implementation-plan.md | 3 ++- internal/loadgen/http.go | 11 ++++++++++- internal/loadgen/http_test.go | 7 ++++++- 8 files changed, 38 insertions(+), 11 deletions(-) diff --git a/README.md b/README.md index 8cfc333..294bdbf 100644 --- a/README.md +++ b/README.md @@ -184,8 +184,19 @@ go run ./cmd/proxy-loadgen ` ``` 使用 `-duration 30s -rate 5000` 可运行限速场景;省略 `-rate` 时固定数量 worker -会饱和发送。命令输出 JSON 报告,包含成功/失败分类、固定内存的延迟分位上界、吞吐和 -Go 运行时内存/GC 快照。它不包含长连接保持或 Extract 场景,也不构成 100,000 QPS 证明。 +会饱和发送。`-method`、重复的 `-header` 与 `-body` 可组合用于 Distribution 的提取接口: + +```powershell +go run ./cmd/proxy-loadgen ` + -target http://CONTROLLER_HOST:8081/api/v1/proxies/extract ` + -method POST -header "Content-Type: application/json" ` + -header "X-API-Key: DISTRIBUTION_API_KEY" -body '{"count":1}' ` + -requests 1000 -concurrency 32 -timeout 10s +``` + +命令输出 JSON 报告,包含成功/失败分类、固定内存的延迟分位上界、吞吐和 Go +运行时内存/GC 快照。它尚不包含长连接保持、Extract 的专用数据准备/结果校验, +也不构成 100,000 QPS 证明。 ## 关键配置与入口 diff --git a/cmd/proxy-loadgen/main.go b/cmd/proxy-loadgen/main.go index 33f40d3..c5798a9 100644 --- a/cmd/proxy-loadgen/main.go +++ b/cmd/proxy-loadgen/main.go @@ -31,6 +31,7 @@ func execute(ctx context.Context, args []string, run loadRun, stdout, stderr io. targetURL := flags.String("target", "", "HTTP or HTTPS target URL") proxyURL := flags.String("proxy", "", "optional HTTP or HTTPS forward proxy URL") method := flags.String("method", http.MethodGet, "HTTP method") + body := flags.String("body", "", "UTF-8 request body") requests := flags.Int("requests", 0, "fixed request count; mutually exclusive with -duration") duration := flags.Duration("duration", 0, "time-boxed workload duration") rate := flags.Int("rate", 0, "maximum request starts per second for -duration; zero saturates workers") @@ -54,7 +55,7 @@ func execute(ctx context.Context, args []string, run loadRun, stdout, stderr io. return 2 } report, err := run(ctx, loadgen.Options{ - TargetURL: *targetURL, ProxyURL: *proxyURL, Method: *method, Headers: parsedHeaders, + TargetURL: *targetURL, ProxyURL: *proxyURL, Method: *method, Headers: parsedHeaders, RequestBody: []byte(*body), Requests: *requests, Duration: *duration, Rate: *rate, Concurrency: *concurrency, RequestTimeout: *timeout, }) if err != nil { diff --git a/cmd/proxy-loadgen/main_test.go b/cmd/proxy-loadgen/main_test.go index acd3391..537354d 100644 --- a/cmd/proxy-loadgen/main_test.go +++ b/cmd/proxy-loadgen/main_test.go @@ -15,14 +15,14 @@ func TestExecutePassesBoundedWorkloadAndWritesJSON(t *testing.T) { var stdout, stderr bytes.Buffer code := execute(context.Background(), []string{ "-target", "https://target.example/path", "-proxy", "http://gateway.example:8080", "-method", "post", - "-requests", "3", "-concurrency", "2", "-timeout", "2s", "-header", "X-Run: fixed", + "-requests", "3", "-concurrency", "2", "-timeout", "2s", "-header", "X-Run: fixed", "-body", `{"count":1}`, }, func(_ context.Context, options loadgen.Options) (loadgen.Report, error) { received = options return loadgen.Report{Requests: 3, Completed: 3, Succeeded: 3, Duration: time.Second}, nil }, &stdout, &stderr) if code != 0 || received.TargetURL != "https://target.example/path" || received.ProxyURL != "http://gateway.example:8080" || received.Method != "post" || received.Requests != 3 || received.Duration != 0 || received.Concurrency != 2 || - received.Headers.Get("X-Run") != "fixed" || stderr.Len() != 0 { + received.Headers.Get("X-Run") != "fixed" || string(received.RequestBody) != `{"count":1}` || stderr.Len() != 0 { t.Fatalf("execute() = %d; options=%+v stderr=%q", code, received, stderr.String()) } var report loadgen.Report diff --git a/docs/design/architecture.md b/docs/design/architecture.md index 52071dd..62f4e3d 100644 --- a/docs/design/architecture.md +++ b/docs/design/architecture.md @@ -79,8 +79,8 @@ Provider、Pool、Routing、Distribution 在首版需要共享事务和一致性 ### 3.4 proxy-loadgen -- 当前实现 HTTP 请求场景:固定请求数或固定时长,受限并发与可选目标 QPS;HTTPS - 目标经 HTTP Gateway 时由 Transport 走 CONNECT。 +- 当前实现 HTTP 请求场景:固定请求数或固定时长,受限并发与可选目标 QPS;可配置方法、 + 重复请求头和请求体。HTTPS 目标经 HTTP Gateway 时由 Transport 走 CONNECT。 - 输出场景、延迟分位上界、错误分类、吞吐和 Go 内存/GC 快照。CONNECT 长连接、 Extract 并发、进程 CPU/RSS/句柄与网络采样仍需补齐。 diff --git a/docs/design/project-structure.md b/docs/design/project-structure.md index f643d5a..5fd1954 100644 --- a/docs/design/project-structure.md +++ b/docs/design/project-structure.md @@ -54,7 +54,7 @@ deadline 内执行 HTTP/HTTPS/SOCKS5 BASIC、EGRESS 和 TARGET 探测并微批 ### proxy-loadgen 负载工具当前生成有界 HTTP 请求,支持经 Gateway 请求 HTTPS 目标、固定请求数或时长、 -目标 QPS、连接复用和 JSON 指标输出。延迟统计使用固定大小直方图,不会因长时间高 QPS +目标 QPS、可重复请求头和请求体、连接复用和 JSON 指标输出。延迟统计使用固定大小直方图,不会因长时间高 QPS 运行积压样本。CONNECT 长连接、Extract 和故障注入场景仍待补齐;它是 100k QPS 结论的 证据工具,不是业务进程。 diff --git a/docs/development/implementation-plan.md b/docs/development/implementation-plan.md index f2fc9de..5189bc4 100644 --- a/docs/development/implementation-plan.md +++ b/docs/development/implementation-plan.md @@ -295,7 +295,8 @@ Controller 即可生效;Redis 任务存储仍仅承载 BASIC,未扩展 EGRES 补充进度(2026-08-02):已新增 `proxy-loadgen` HTTP 场景。固定请求数和固定时长两种 模式均通过固定 worker 数与有界派发通道执行,可选 QPS 限速;报告使用固定大小延迟直方图, -输出状态分类、吞吐和 Go 内存/GC 快照。CONNECT 长连接、Extract、故障注入以及代表性集群 +输出状态分类、吞吐和 Go 内存/GC 快照。通用 `method/header/body` 参数可覆盖 Distribution +提取 HTTP 请求;CONNECT 长连接、Extract 的专用数据准备与结果校验、故障注入以及代表性集群 报告仍未实现。 ## Task 12: Machine-readable Contracts diff --git a/internal/loadgen/http.go b/internal/loadgen/http.go index 443d323..29bae50 100644 --- a/internal/loadgen/http.go +++ b/internal/loadgen/http.go @@ -3,6 +3,7 @@ package loadgen import ( + "bytes" "context" "errors" "io" @@ -27,6 +28,7 @@ type Options struct { ProxyURL string Method string Headers http.Header + RequestBody []byte Requests int Duration time.Duration Rate int @@ -110,6 +112,8 @@ func Run(ctx context.Context, options Options) (Report, error) { func newClient(options Options) (*http.Client, Options, error) { normalized := options normalized.Method = strings.ToUpper(strings.TrimSpace(normalized.Method)) + normalized.Headers = normalized.Headers.Clone() + normalized.RequestBody = bytes.Clone(normalized.RequestBody) if normalized.Method == "" { normalized.Method = http.MethodGet } @@ -245,7 +249,12 @@ func executeOne(ctx context.Context, client *http.Client, options Options, stats started := time.Now() requestContext, cancel := context.WithTimeout(ctx, options.RequestTimeout) defer cancel() - request, err := http.NewRequestWithContext(requestContext, options.Method, options.TargetURL, nil) + request, err := http.NewRequestWithContext( + requestContext, + options.Method, + options.TargetURL, + bytes.NewReader(options.RequestBody), + ) if err != nil { stats.failed.Add(1) stats.requestErrors.Add(1) diff --git a/internal/loadgen/http_test.go b/internal/loadgen/http_test.go index fb3a318..64ef3b1 100644 --- a/internal/loadgen/http_test.go +++ b/internal/loadgen/http_test.go @@ -3,6 +3,7 @@ package loadgen import ( "context" "errors" + "io" "net/http" "net/http/httptest" "sync/atomic" @@ -16,13 +17,17 @@ func TestRunExecutesFixedBoundedHTTPWorkload(t *testing.T) { if request.Method != http.MethodPost || request.Header.Get("X-Scenario") != "fixed" { t.Errorf("request = %s %q", request.Method, request.Header.Get("X-Scenario")) } + body, err := io.ReadAll(request.Body) + if err != nil || string(body) != `{"count":1}` { + t.Errorf("request body = %q, error = %v", body, err) + } requests.Add(1) response.WriteHeader(http.StatusNoContent) })) defer server.Close() report, err := Run(context.Background(), Options{ - TargetURL: server.URL, Method: http.MethodPost, Headers: http.Header{"X-Scenario": {"fixed"}}, + TargetURL: server.URL, Method: http.MethodPost, Headers: http.Header{"X-Scenario": {"fixed"}}, RequestBody: []byte(`{"count":1}`), Requests: 12, Concurrency: 3, RequestTimeout: time.Second, }) if err != nil || report.Requests != 12 || report.Completed != 12 || report.Succeeded != 12 || report.Failed != 0 ||