From c55920cade089843d283b5c5c91ea7b77a55ba1a Mon Sep 17 00:00:00 2001 From: tonghzhang <949674415@qq.com> Date: Mon, 27 Jul 2026 16:31:59 +0800 Subject: [PATCH 1/9] test(perf): add hot-path diagnostic metrics --- test/threehost/README.md | 75 +- test/threehost/collect-compose-stats.sh | 231 ++++- test/threehost/collect-prometheus.sh | 180 +++- test/threehost/render_html.py | 256 +++++ test/threehost/report-template.html | 1148 +++++++++++++++++++++++ test/threehost/run-complete-client.sh | 10 +- test/threehost/run-diagnostic-client.sh | 45 + test/threehost/summarize.py | 630 ++++++++++++- 8 files changed, 2518 insertions(+), 57 deletions(-) create mode 100644 test/threehost/render_html.py create mode 100644 test/threehost/report-template.html create mode 100644 test/threehost/run-diagnostic-client.sh diff --git a/test/threehost/README.md b/test/threehost/README.md index 17c5e49..b34a1a8 100644 --- a/test/threehost/README.md +++ b/test/threehost/README.md @@ -154,6 +154,26 @@ bash test/threehost/run-client.sh 结果写入 `test-results/threehost//`。每个 case 都有 k6 summary JSON 和文本日志; `SAVE_RAW_METRICS=true` 时还会保存逐指标 JSON,文件会明显变大。 +### 第一轮热路径诊断 + +先不要运行完整矩阵。准备好 `benchmark.env` 后执行约 19 分钟的固定诊断轮: + +```bash +cp test/threehost/benchmark.env.example test/threehost/benchmark.env +editor test/threehost/benchmark.env +bash test/threehost/run-diagnostic-client.sh +``` + +该入口复用完整 runner,只运行 smoke、100 RPS 预热、`500/750/1000/1500 RPS × 2m` +和 `500 RPS × 10m` 耐久。它不会修改网关 Queue、Redis、PostgreSQL 或限流参数。客户端 +最大 VU 提高到 4096,避免把 1500 RPS 过载点过早归因于默认 2048 VU 上限。 + +诊断轮必须与网关机的 Prometheus 和容器采集同时运行;两台服务器的 +`DURATION_SECONDS=1800` 即可覆盖测试和余量。最终 HTML 的“热路径诊断”可以逐 case +切换,显示认证、Usage Begin、授权、Redis 限流、Quota、Provider Queue/调用、Usage +Finalize 的 P50/P95/P99,以及 PostgreSQL/Redis pool wait、Go CPU/RSS/heap/goroutine/GC +和稳定内部错误码。缺少重叠监控时这些字段不会被当成 0,而会给出证据警告。 + ### 完整性能与故障套件 `run-client.sh` 适合先验证三机连通性。正式测试改用完整配置: @@ -215,11 +235,44 @@ PRE_ALLOCATED_VUS >= RATE × 直连 p95 秒数 × 1.2 在客户端开始前,在另外两台服务器的终端中运行采集器。`DURATION_SECONDS` 应覆盖预热、 所有 repetitions、冷却间隔和可靠性测试。 +采集器现在每 10 秒输出一次进度,例如: + +```text +[2026-07-26T15:20:00Z] docker-stats state=collecting elapsed=00:10:00 remaining=02:20:00 snapshots=297 records=1188 errors=0 last=2026-07-26T15:19:59.842Z +``` + +`state=collecting` 表示正在采样,`remaining` 是离计划结束的剩余时间,`last` 是最后一个成功 +样本。对应的 `*-status.json` 会在每轮采样后原子更新;正常结束记为 `completed`,收到 +Ctrl+C、TERM 或 HUP 记为 `interrupted`,其他异常记为 `failed` 并保留原因。用 +`PROGRESS_SECONDS` 可以调整终端刷新间隔。 + +如果状态仍是 `collecting`,但 `updated_at` 已超过两个实际采样周期没有变化,应视为 +采集器已经失联或卡住,不要继续启动正式压测。可以在另一个终端持续查看状态: + +```bash +watch -n 2 'python3 -m json.tool test-results/threehost//gateway-stats-status.json' +``` + +不要让采集器依附于可能关闭的普通 SSH 会话。推荐先进入一个持久的 `tmux` 会话,再执行 +下面的前台命令;这样既能直接看到进度,又能在 SSH 断开后继续采集: + +```bash +tmux new -s model-velo-monitor +# 在 tmux 中执行本机对应的采集命令。 +# 按 Ctrl+B,再按 D,退出但不停止采集。 +tmux attach -t model-velo-monitor +``` + +如果一次客户端运行 smoke 失败或被中断,不要停止采集器后沿用同一批文件假装正式运行。 +修好问题后应使用新的 `RUN_ID`(或新的 attempt 目录),确认三台监控仍显示 +`state=collecting`,再启动正式客户端。 + 网关机: ```bash RUN_ID=20260725T120000Z DURATION_SECONDS=9000 \ +PROGRESS_SECONDS=10 \ OUTPUT_FILE="test-results/threehost/$RUN_ID/gateway-stats.jsonl" \ bash test/threehost/collect-compose-stats.sh ``` @@ -229,6 +282,7 @@ bash test/threehost/collect-compose-stats.sh ```bash RUN_ID=20260725T120000Z DURATION_SECONDS=9000 \ +PROGRESS_SECONDS=10 \ OUTPUT_DIR="test-results/threehost/$RUN_ID/prometheus" \ bash test/threehost/collect-prometheus.sh ``` @@ -238,6 +292,7 @@ bash test/threehost/collect-prometheus.sh ```bash RUN_ID=20260725T120000Z DURATION_SECONDS=9000 \ +PROGRESS_SECONDS=10 \ COMPOSE_FILE=test/threehost/upstream.compose.yaml \ SERVICES=main,fail,fallback \ OUTPUT_FILE="test-results/threehost/$RUN_ID/upstream-stats.jsonl" \ @@ -247,6 +302,10 @@ bash test/threehost/collect-compose-stats.sh 采集器按 `INTERVAL_SECONDS` 使用一次 `docker stats --no-stream` 批量记录所有目标容器的 CPU、内存、网络和块 I/O,同时单独记录服务到容器 ID 的映射、镜像 ID、Docker 版本和 宿主机信息;它不会执行 `docker compose config`,避免把 `.env` 密钥写入结果。 +`INTERVAL_SECONDS` 是最小采样周期:如果一次 `docker stats` 本身超过该时间,下一轮会 +立即开始,不会再额外 sleep;因此 1 秒配置在当前 Docker 环境中通常得到约 2 秒而不是 +原先的约 3 秒实际间隔。Docker 短暂失败时会刷新容器 ID 并继续;默认连续失败 30 次才 +将监控标记为 `failed`。 客户端套件结束后,在网关机保存一次可直接诊断的最终证据: @@ -277,6 +336,8 @@ scp -r gateway-host:"~/model-velo/$RESULT_DIR/gateway-stats.jsonl" \ "$RESULT_DIR/" scp -r gateway-host:"~/model-velo/$RESULT_DIR/gateway-stats-metadata.txt" \ "$RESULT_DIR/" +scp -r gateway-host:"~/model-velo/$RESULT_DIR/gateway-stats-status.json" \ + "$RESULT_DIR/" scp -r gateway-host:"~/model-velo/$RESULT_DIR/prometheus" \ "$RESULT_DIR/" scp -r gateway-host:"~/model-velo/$RESULT_DIR/gateway-evidence" \ @@ -285,16 +346,19 @@ scp -r upstream-host:"~/model-velo/$RESULT_DIR/upstream-stats.jsonl" \ "$RESULT_DIR/" scp -r upstream-host:"~/model-velo/$RESULT_DIR/upstream-stats-metadata.txt" \ "$RESULT_DIR/" +scp -r upstream-host:"~/model-velo/$RESULT_DIR/upstream-stats-status.json" \ + "$RESULT_DIR/" python3 test/threehost/summarize.py "$RESULT_DIR" tar -czf "$RUN_ID.tar.gz" -C test-results/threehost "$RUN_ID" ``` -把最后的 `.tar.gz` 给分析者即可。`summary.md` 是人读摘要,`summary.json` 保留 -全部机器可读 case、分位数、状态计数、直连差值、逐 Chunk SSE、资源、Prometheus 和 -Usage 证据。原始 `*.log`、`*-summary.json`、`*-stream.json`、`*-upstream.json` 和 -带时间戳采样也都保留;下一轮可以据此定位是客户端 VU、网关 CPU/Queue/Redis、Usage -Worker、上游放大还是长尾问题,再改对应代码。 +把最后的 `.tar.gz` 给分析者即可。`summary.html` 是不依赖 CDN、可直接在浏览器 +打开的交互报告,`summary.md` 是终端友好的摘要,`summary.json` 保留全部机器可读 case、 +分位数、状态计数、直连差值、逐 Chunk SSE、资源、Prometheus 和 Usage 证据。原始 +`*.log`、`*-summary.json`、`*-stream.json`、`*-upstream.json` 和带时间戳采样也都保留; +下一轮可以据此定位是客户端 VU、网关 CPU/Queue/Redis、Usage Worker、上游放大还是长尾 +问题,再改对应代码。 汇总还会把正常网关 case 的实际 HTTP 请求数和落库 Usage Event 数对账;smoke、 reliability 和独立限流用例因包含非 Chat 请求或前置拒绝,不纳入该等式。 @@ -350,6 +414,7 @@ reliability 和独立限流用例因包含非 Chat 请求或前置拒绝,不 | `prepare-gateway-env.sh` | 生成不入库的三机网关 `.env` | | `run-client.sh` | 固定顺序运行并保存 k6 结果 | | `run-complete-client.sh` | 执行完整矩阵、保存 case 时间与上游计数 | +| `run-diagnostic-client.sh` | 执行约 19 分钟的第一轮热路径诊断 | | `collect-host-stats.sh` | 自动采集客户端 CPU、内存、Load 和负载进程 RSS | | `collect-compose-stats.sh` | 在被测主机本地采集容器资源 | | `collect-prometheus.sh` | 连续采集网关和 Worker 的低基数指标 | diff --git a/test/threehost/collect-compose-stats.sh b/test/threehost/collect-compose-stats.sh index 0508220..51077a6 100644 --- a/test/threehost/collect-compose-stats.sh +++ b/test/threehost/collect-compose-stats.sh @@ -9,6 +9,9 @@ services_csv=${SERVICES:-gateway,usage-worker,postgres,redis} duration_seconds=${DURATION_SECONDS:-600} interval_seconds=${INTERVAL_SECONDS:-1} output_file=${OUTPUT_FILE:-"$repo_root/test-results/threehost/compose-stats.jsonl"} +progress_seconds=${PROGRESS_SECONDS:-10} +max_consecutive_errors=${MAX_CONSECUTIVE_ERRORS:-30} +status_file=${STATUS_FILE:-"${output_file%.jsonl}-status.json"} for command_name in docker date git; do if ! command -v "$command_name" >/dev/null 2>&1; then @@ -28,42 +31,175 @@ if [[ ! "$interval_seconds" =~ ^[1-9][0-9]*$ ]]; then printf 'INTERVAL_SECONDS must be a positive integer\n' >&2 exit 1 fi +if [[ ! "$progress_seconds" =~ ^[1-9][0-9]*$ ]]; then + printf 'PROGRESS_SECONDS must be a positive integer\n' >&2 + exit 1 +fi +if [[ ! "$max_consecutive_errors" =~ ^[1-9][0-9]*$ ]]; then + printf 'MAX_CONSECUTIVE_ERRORS must be a positive integer\n' >&2 + exit 1 +fi IFS=',' read -r -a services <<<"$services_csv" -declare -a container_ids=() declare -a service_names=() for raw_service in "${services[@]}"; do service=$(printf '%s' "$raw_service" | tr -d '[:space:]') if [[ -z "$service" ]]; then continue fi - container_id=$(docker compose -f "$compose_file" ps -q "$service") - if [[ -z "$container_id" ]]; then - printf 'service is not running: %s\n' "$service" >&2 - exit 1 - fi service_names+=("$service") - container_ids+=("$container_id") done -if ((${#container_ids[@]} == 0)); then +if ((${#service_names[@]} == 0)); then printf 'SERVICES did not contain any service names\n' >&2 exit 1 fi -mkdir -p "$(dirname -- "$output_file")" +declare -a container_ids=() +refresh_container_ids() { + local service + local container_id + local -a refreshed_ids=() + for service in "${service_names[@]}"; do + container_id=$(docker compose -f "$compose_file" ps -q "$service") || return 1 + if [[ -z "$container_id" ]]; then + return 1 + fi + refreshed_ids+=("$container_id") + done + container_ids=("${refreshed_ids[@]}") +} +if ! refresh_container_ids; then + printf 'one or more services are not running: %s\n' "$services_csv" >&2 + exit 1 +fi + +mkdir -p "$(dirname -- "$output_file")" "$(dirname -- "$status_file")" metadata_file="${output_file%.jsonl}-metadata.txt" -if [[ -e "$output_file" || -e "$metadata_file" ]]; then - printf 'refusing to overwrite existing stats output: %s or %s\n' \ +if [[ -e "$output_file" || -e "$metadata_file" || -e "$status_file" ]]; then + printf 'refusing to overwrite existing stats output: %s, %s, or %s\n' \ "$output_file" \ - "$metadata_file" >&2 + "$metadata_file" \ + "$status_file" >&2 exit 1 fi + +format_duration() { + local total_seconds=$1 + printf '%02d:%02d:%02d' \ + "$((total_seconds / 3600))" \ + "$(((total_seconds % 3600) / 60))" \ + "$((total_seconds % 60))" +} + +json_escape() { + local value=$1 + value=${value//\\/\\\\} + value=${value//\"/\\\"} + value=${value//$'\n'/\\n} + value=${value//$'\r'/\\r} + value=${value//$'\t'/\\t} + printf '%s' "$value" +} + +monitor_started=$SECONDS +started_at=$(date -u +%Y-%m-%dT%H:%M:%SZ) +deadline=$((monitor_started + duration_seconds)) +next_progress=$monitor_started +snapshots=0 +records=0 +errors=0 +consecutive_errors=0 +last_sample_at= +last_error= +stop_requested=false +termination_signal= +termination_exit_code=1 +completed=false + +write_status() { + local state=$1 + local reason=${2:-} + local elapsed=$((SECONDS - monitor_started)) + local remaining=$((deadline - SECONDS)) + local temporary_status="${status_file}.tmp.$$" + if ((remaining < 0)); then + remaining=0 + fi + printf '{"state":"%s","pid":%d,"started_at":"%s","updated_at":"%s",' \ + "$state" \ + "$$" \ + "$started_at" \ + "$(date -u +%Y-%m-%dT%H:%M:%SZ)" >"$temporary_status" + printf '"duration_seconds":%d,"elapsed_seconds":%d,"remaining_seconds":%d,' \ + "$duration_seconds" \ + "$elapsed" \ + "$remaining" >>"$temporary_status" + printf '"snapshots":%d,"records":%d,"errors":%d,"consecutive_errors":%d,' \ + "$snapshots" \ + "$records" \ + "$errors" \ + "$consecutive_errors" >>"$temporary_status" + printf '"last_sample_at":"%s","reason":"%s","last_error":"%s"}\n' \ + "$last_sample_at" \ + "$(json_escape "$reason")" \ + "$(json_escape "$last_error")" >>"$temporary_status" + mv -- "$temporary_status" "$status_file" +} + +print_progress() { + local state=$1 + local elapsed=$((SECONDS - monitor_started)) + local remaining=$((deadline - SECONDS)) + if ((remaining < 0)); then + remaining=0 + fi + printf '[%s] docker-stats state=%s elapsed=%s remaining=%s snapshots=%d records=%d errors=%d last=%s\n' \ + "$(date -u +%Y-%m-%dT%H:%M:%SZ)" \ + "$state" \ + "$(format_duration "$elapsed")" \ + "$(format_duration "$remaining")" \ + "$snapshots" \ + "$records" \ + "$errors" \ + "${last_sample_at:-none}" +} + +handle_signal() { + termination_signal=$1 + termination_exit_code=$2 + stop_requested=true +} + +finish_monitor() { + local exit_code=$1 + local state=failed + local reason="exit:$exit_code" + set +e + if [[ "$completed" == "true" ]]; then + state=completed + reason=duration_reached + elif [[ -n "$termination_signal" ]]; then + state=interrupted + reason="signal:$termination_signal" + fi + write_status "$state" "$reason" + print_progress "$state" +} + +trap 'handle_signal INT 130' INT +trap 'handle_signal TERM 143' TERM +trap 'handle_signal HUP 129' HUP +trap 'finish_monitor $?' EXIT + { - printf 'captured_at=%s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" + printf 'captured_at=%s\n' "$started_at" printf 'compose_file=%s\n' "$compose_file" printf 'services=%s\n' "$services_csv" printf 'duration_seconds=%s\n' "$duration_seconds" printf 'interval_seconds=%s\n' "$interval_seconds" + printf 'progress_seconds=%s\n' "$progress_seconds" + printf 'max_consecutive_errors=%s\n' "$max_consecutive_errors" + printf 'status_file=%s\n' "$status_file" printf 'host=%s\n' "$(uname -a)" printf 'commit=%s\n' "$(git -C "$repo_root" rev-parse HEAD)" if [[ -n "$(git -C "$repo_root" status --porcelain)" ]]; then @@ -95,23 +231,66 @@ fi fi } >"$metadata_file" -deadline=$((SECONDS + duration_seconds)) : >"$output_file" -while ((SECONDS < deadline)); do +write_status starting +print_progress starting +while ((SECONDS < deadline)) && [[ "$stop_requested" == "false" ]]; do + iteration_started=$SECONDS timestamp=$(date -u +%Y-%m-%dT%H:%M:%S.%3NZ) - while IFS= read -r stats; do - if [[ -n "$stats" ]]; then - printf '{"timestamp":"%s","stats":%s}\n' \ - "$timestamp" \ - "$stats" >>"$output_file" - fi - done < <( - docker stats \ + stats_output= + if stats_output=$(docker stats \ --no-stream \ --format '{{json .}}' \ - "${container_ids[@]}" - ) - sleep "$interval_seconds" + "${container_ids[@]}" 2>&1); then + snapshot_records=0 + while IFS= read -r stats; do + if [[ -n "$stats" ]]; then + printf '{"timestamp":"%s","stats":%s}\n' \ + "$timestamp" \ + "$stats" >>"$output_file" + snapshot_records=$((snapshot_records + 1)) + records=$((records + 1)) + fi + done <<<"$stats_output" + if ((snapshot_records > 0)); then + snapshots=$((snapshots + 1)) + consecutive_errors=0 + last_error= + last_sample_at=$timestamp + else + errors=$((errors + 1)) + consecutive_errors=$((consecutive_errors + 1)) + last_error="docker stats returned no records" + fi + else + errors=$((errors + 1)) + consecutive_errors=$((consecutive_errors + 1)) + last_error=$stats_output + refresh_container_ids || true + fi + + write_status collecting + if ((SECONDS >= next_progress)); then + print_progress collecting + next_progress=$((SECONDS + progress_seconds)) + fi + if ((consecutive_errors >= max_consecutive_errors)); then + printf 'docker stats failed %d consecutive times: %s\n' \ + "$consecutive_errors" \ + "$last_error" >&2 + exit 1 + fi + + iteration_elapsed=$((SECONDS - iteration_started)) + sleep_seconds=$((interval_seconds - iteration_elapsed)) + if ((sleep_seconds > 0)) && ((SECONDS < deadline)); then + sleep "$sleep_seconds" & + wait $! || true + fi done +if [[ "$stop_requested" == "true" ]]; then + exit "$termination_exit_code" +fi +completed=true printf 'wrote %s and %s\n' "$output_file" "$metadata_file" diff --git a/test/threehost/collect-prometheus.sh b/test/threehost/collect-prometheus.sh index 2f7a3e2..f35ac5f 100644 --- a/test/threehost/collect-prometheus.sh +++ b/test/threehost/collect-prometheus.sh @@ -24,6 +24,9 @@ set +a duration_seconds=${DURATION_SECONDS:-7200} interval_seconds=${INTERVAL_SECONDS:-1} output_dir=${OUTPUT_DIR:-"$repo_root/test-results/threehost/metrics"} +progress_seconds=${PROGRESS_SECONDS:-10} +max_consecutive_errors=${MAX_CONSECUTIVE_ERRORS:-30} +status_file=${STATUS_FILE:-"$output_dir/monitor-status.json"} bind_address=${MODEL_VELO_HTTP_BIND:-127.0.0.1} if [[ "$bind_address" == "0.0.0.0" || "$bind_address" == "[::]" ]]; then bind_address=127.0.0.1 @@ -40,6 +43,14 @@ if [[ ! "$interval_seconds" =~ ^[1-9][0-9]*$ ]]; then printf 'INTERVAL_SECONDS must be a positive integer\n' >&2 exit 1 fi +if [[ ! "$progress_seconds" =~ ^[1-9][0-9]*$ ]]; then + printf 'PROGRESS_SECONDS must be a positive integer\n' >&2 + exit 1 +fi +if [[ ! "$max_consecutive_errors" =~ ^[1-9][0-9]*$ ]]; then + printf 'MAX_CONSECUTIVE_ERRORS must be a positive integer\n' >&2 + exit 1 +fi if [[ -e "$output_dir" ]]; then printf 'refusing to overwrite metrics output directory: %s\n' "$output_dir" >&2 exit 1 @@ -52,10 +63,121 @@ metadata_output="$output_dir/prometheus-metadata.txt" : >"$gateway_output" : >"$worker_output" +format_duration() { + local total_seconds=$1 + printf '%02d:%02d:%02d' \ + "$((total_seconds / 3600))" \ + "$(((total_seconds % 3600) / 60))" \ + "$((total_seconds % 60))" +} + +json_escape() { + local value=$1 + value=${value//\\/\\\\} + value=${value//\"/\\\"} + value=${value//$'\n'/\\n} + value=${value//$'\r'/\\r} + value=${value//$'\t'/\\t} + printf '%s' "$value" +} + +monitor_started=$SECONDS +started_at=$(date -u +%Y-%m-%dT%H:%M:%SZ) +deadline=$((monitor_started + duration_seconds)) +next_progress=$monitor_started +snapshots=0 +scrapes=0 +errors=0 +consecutive_errors=0 +last_sample_at= +last_error= +stop_requested=false +termination_signal= +termination_exit_code=1 +completed=false + +write_status() { + local state=$1 + local reason=${2:-} + local elapsed=$((SECONDS - monitor_started)) + local remaining=$((deadline - SECONDS)) + local temporary_status="${status_file}.tmp.$$" + if ((remaining < 0)); then + remaining=0 + fi + printf '{"state":"%s","pid":%d,"started_at":"%s","updated_at":"%s",' \ + "$state" \ + "$$" \ + "$started_at" \ + "$(date -u +%Y-%m-%dT%H:%M:%SZ)" >"$temporary_status" + printf '"duration_seconds":%d,"elapsed_seconds":%d,"remaining_seconds":%d,' \ + "$duration_seconds" \ + "$elapsed" \ + "$remaining" >>"$temporary_status" + printf '"snapshots":%d,"scrapes":%d,"errors":%d,"consecutive_errors":%d,' \ + "$snapshots" \ + "$scrapes" \ + "$errors" \ + "$consecutive_errors" >>"$temporary_status" + printf '"last_sample_at":"%s","reason":"%s","last_error":"%s"}\n' \ + "$last_sample_at" \ + "$(json_escape "$reason")" \ + "$(json_escape "$last_error")" >>"$temporary_status" + mv -- "$temporary_status" "$status_file" +} + +print_progress() { + local state=$1 + local elapsed=$((SECONDS - monitor_started)) + local remaining=$((deadline - SECONDS)) + if ((remaining < 0)); then + remaining=0 + fi + printf '[%s] prometheus state=%s elapsed=%s remaining=%s snapshots=%d scrapes=%d errors=%d last=%s\n' \ + "$(date -u +%Y-%m-%dT%H:%M:%SZ)" \ + "$state" \ + "$(format_duration "$elapsed")" \ + "$(format_duration "$remaining")" \ + "$snapshots" \ + "$scrapes" \ + "$errors" \ + "${last_sample_at:-none}" +} + +handle_signal() { + termination_signal=$1 + termination_exit_code=$2 + stop_requested=true +} + +finish_monitor() { + local exit_code=$1 + local state=failed + local reason="exit:$exit_code" + set +e + if [[ "$completed" == "true" ]]; then + state=completed + reason=duration_reached + elif [[ -n "$termination_signal" ]]; then + state=interrupted + reason="signal:$termination_signal" + fi + write_status "$state" "$reason" + print_progress "$state" +} + +trap 'handle_signal INT 130' INT +trap 'handle_signal TERM 143' TERM +trap 'handle_signal HUP 129' HUP +trap 'finish_monitor $?' EXIT + { - printf 'captured_at=%s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" + printf 'captured_at=%s\n' "$started_at" printf 'duration_seconds=%s\n' "$duration_seconds" printf 'interval_seconds=%s\n' "$interval_seconds" + printf 'progress_seconds=%s\n' "$progress_seconds" + printf 'max_consecutive_errors=%s\n' "$max_consecutive_errors" + printf 'status_file=%s\n' "$status_file" printf 'gateway_url=%s\n' "$gateway_url" printf 'worker_url=%s\n' "$worker_url" printf 'commit=%s\n' "$(git -C "$repo_root" rev-parse HEAD)" @@ -75,18 +197,60 @@ capture() { printf '# snapshot %s\n' "$timestamp" >>"$output" if ! payload=$(curl "${curl_args[@]}" "$url"); then printf '# scrape_error\n' >>"$output" - return + last_error="scrape failed: $url" + return 1 fi printf '%s\n' "$payload" | - awk '/^model_velo_[A-Za-z0-9_:]+({[^}]*})? [^ ]+/' >>"$output" + awk '/^(model_velo_|go_|process_)[A-Za-z0-9_:]+({[^}]*})? [^ ]+/' >>"$output" + return 0 } -deadline=$((SECONDS + duration_seconds)) -while ((SECONDS < deadline)); do +write_status starting +print_progress starting +while ((SECONDS < deadline)) && [[ "$stop_requested" == "false" ]]; do + iteration_started=$SECONDS timestamp=$(date -u +%Y-%m-%dT%H:%M:%S.%3NZ) - capture "$gateway_url" "$gateway_output" "$timestamp" - capture "$worker_url" "$worker_output" "$timestamp" - sleep "$interval_seconds" + snapshot_errors=0 + if ! capture "$gateway_url" "$gateway_output" "$timestamp"; then + snapshot_errors=$((snapshot_errors + 1)) + fi + if ! capture "$worker_url" "$worker_output" "$timestamp"; then + snapshot_errors=$((snapshot_errors + 1)) + fi + snapshots=$((snapshots + 1)) + scrapes=$((scrapes + 2)) + errors=$((errors + snapshot_errors)) + if ((snapshot_errors == 2)); then + consecutive_errors=$((consecutive_errors + 1)) + else + last_sample_at=$timestamp + consecutive_errors=0 + if ((snapshot_errors == 0)); then + last_error= + fi + fi + + write_status collecting + if ((SECONDS >= next_progress)); then + print_progress collecting + next_progress=$((SECONDS + progress_seconds)) + fi + if ((consecutive_errors >= max_consecutive_errors)); then + printf 'both Prometheus scrapes failed %d consecutive times\n' \ + "$consecutive_errors" >&2 + exit 1 + fi + + iteration_elapsed=$((SECONDS - iteration_started)) + sleep_seconds=$((interval_seconds - iteration_elapsed)) + if ((sleep_seconds > 0)) && ((SECONDS < deadline)); then + sleep "$sleep_seconds" & + wait $! || true + fi done +if [[ "$stop_requested" == "true" ]]; then + exit "$termination_exit_code" +fi +completed=true printf 'wrote metrics under %s\n' "$output_dir" diff --git a/test/threehost/render_html.py b/test/threehost/render_html.py new file mode 100644 index 0000000..2b31ec7 --- /dev/null +++ b/test/threehost/render_html.py @@ -0,0 +1,256 @@ +#!/usr/bin/env python3 + +import json +import re +import sys +from datetime import datetime, timezone +from pathlib import Path + + +PLACEHOLDER = "__MODEL_VELO_REPORT_DATA__" + + +def read_text(path): + try: + return path.read_text(encoding="utf-8", errors="replace") + except OSError: + return "" + + +def read_json(path): + try: + return json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + return {} + + +def parse_metadata(text): + values = {} + for line in text.splitlines(): + if "=" in line: + key, value = line.split("=", 1) + values[key.strip()] = value.strip() + return values + + +def read_json_lines(path): + rows = [] + for line in read_text(path).splitlines(): + try: + rows.append(json.loads(line)) + except json.JSONDecodeError: + continue + return rows + + +def attempt_status(path, current): + if path == current: + return "complete" + name = path.name.lower() + if "smoke-failed" in name: + return "smoke failed" + if "interrupted" in name: + return "interrupted" + if "failed" in name: + return "earlier failed run" + return "other" + + +def collect_attempts(result_dir): + attempts = [] + for path in sorted(item for item in result_dir.parent.iterdir() if item.is_dir()): + files = [item for item in path.rglob("*") if item.is_file()] + metadata_text = read_text(path / "client-metadata.txt") + metadata = parse_metadata(metadata_text) + cases_path = path / "cases.tsv" + case_count = 0 + if cases_path.exists(): + case_count = max(0, len(read_text(cases_path).splitlines()) - 1) + attempts.append( + { + "name": path.name, + "status": attempt_status(path, result_dir), + "case_count": case_count, + "file_count": len(files), + "bytes": sum(item.stat().st_size for item in files), + "commit": metadata.get("commit", ""), + "started_at": metadata.get("started_at", ""), + "ended_at": metadata.get("ended_at", ""), + "has_summary": (path / "summary.json").exists(), + } + ) + return attempts + + +def collect_artifacts(result_dir): + artifacts = [] + for path in sorted( + item + for item in result_dir.rglob("*") + if item.is_file() and item.name != "summary.html" + ): + relative = path.relative_to(result_dir).as_posix() + artifacts.append( + { + "path": relative, + "bytes": path.stat().st_size, + "modified_at": datetime.fromtimestamp( + path.stat().st_mtime, timezone.utc + ).isoformat(), + } + ) + return artifacts + + +def reliability_checks(result_dir): + checks = [] + for line in read_text(result_dir / "reliability.log").splitlines(): + match = re.match(r"^\s*✓\s+(.+?)\s*$", line) + if match and "rate==" not in match.group(1): + checks.append(match.group(1)) + return checks + + +def parse_iso(value): + try: + return datetime.fromisoformat(str(value).replace("Z", "+00:00")) + except (TypeError, ValueError): + return None + + +def capture_relation(first, last, run_start, run_end): + if not all((first, last, run_start, run_end)): + return "unknown" + if last < run_start: + return "before" + if first > run_end: + return "after" + return "overlap" + + +def parse_size(value): + match = re.match(r"^\s*([0-9.]+)\s*([KMGT]?i?B)\s*$", str(value)) + if not match: + return 0.0 + amount = float(match.group(1)) + unit = match.group(2) + powers = {"B": 0, "KB": 1, "KiB": 1, "MB": 2, "MiB": 2, + "GB": 3, "GiB": 3, "TB": 4, "TiB": 4} + base = 1024 if "iB" in unit else 1000 + return amount * base ** powers[unit] + + +def collect_server_captures(result_dir, metadata): + run_start = parse_iso(metadata.get("started_at")) + run_end = parse_iso(metadata.get("ended_at")) + captures = [] + for path in sorted(result_dir.rglob("*stats*.jsonl")): + if path.name == "client-stats.jsonl": + continue + timestamps = [] + containers = {} + for payload in read_json_lines(path): + timestamp = parse_iso(payload.get("timestamp")) + if timestamp: + timestamps.append(timestamp) + stats = payload.get("stats", {}) + name = stats.get("Name") or stats.get("Container") or "unknown" + aggregate = containers.setdefault( + name, + {"samples": 0, "cpu_sum_pct": 0.0, "cpu_max_pct": 0.0, + "memory_max_mb": 0.0}, + ) + cpu = float(str(stats.get("CPUPerc", "0")).rstrip("%") or 0) + memory = str(stats.get("MemUsage", "")).split("/", 1)[0] + aggregate["samples"] += 1 + aggregate["cpu_sum_pct"] += cpu + aggregate["cpu_max_pct"] = max(aggregate["cpu_max_pct"], cpu) + aggregate["memory_max_mb"] = max( + aggregate["memory_max_mb"], parse_size(memory) / (1024 * 1024) + ) + first = min(timestamps) if timestamps else None + last = max(timestamps) if timestamps else None + for aggregate in containers.values(): + aggregate["cpu_avg_pct"] = ( + aggregate.pop("cpu_sum_pct") / max(1, aggregate["samples"]) + ) + captures.append( + { + "path": path.relative_to(result_dir).as_posix(), + "kind": "docker-stats", + "first": first.isoformat() if first else "", + "last": last.isoformat() if last else "", + "relation": capture_relation(first, last, run_start, run_end), + "samples": len(timestamps), + "containers": containers, + } + ) + + for path in sorted(result_dir.rglob("*.promlog")): + timestamps = [ + parse_iso(match.group(1)) + for match in re.finditer( + r"^# snapshot\s+(\S+)\s*$", + read_text(path), + flags=re.MULTILINE, + ) + ] + timestamps = [timestamp for timestamp in timestamps if timestamp] + first = min(timestamps) if timestamps else None + last = max(timestamps) if timestamps else None + captures.append( + { + "path": path.relative_to(result_dir).as_posix(), + "kind": "prometheus", + "first": first.isoformat() if first else "", + "last": last.isoformat() if last else "", + "relation": capture_relation(first, last, run_start, run_end), + "samples": len(timestamps), + "containers": {}, + } + ) + return captures + + +def render(result_dir): + summary_path = result_dir / "summary.json" + if not summary_path.exists(): + raise SystemExit(f"summary not found: {summary_path}") + + template_path = Path(__file__).with_name("report-template.html") + template = read_text(template_path) + if PLACEHOLDER not in template: + raise SystemExit(f"report placeholder not found in {template_path}") + + metadata_text = read_text(result_dir / "client-metadata.txt") + metadata = parse_metadata(metadata_text) + report_data = { + "summary": read_json(summary_path), + "metadata": metadata, + "metadata_text": metadata_text, + "client_stats": read_json_lines(result_dir / "client-stats.jsonl"), + "attempts": collect_attempts(result_dir), + "artifacts": collect_artifacts(result_dir), + "reliability_checks": reliability_checks(result_dir), + "server_captures": collect_server_captures(result_dir, metadata), + "generated_at": datetime.now(timezone.utc).isoformat(), + } + encoded = json.dumps( + report_data, ensure_ascii=False, separators=(",", ":") + ).replace("<", "\\u003c") + output = result_dir / "summary.html" + output.write_text(template.replace(PLACEHOLDER, encoded), encoding="utf-8") + print(f"wrote {output}") + + +def main(): + if len(sys.argv) != 2: + raise SystemExit(f"usage: {Path(sys.argv[0]).name} ") + result_dir = Path(sys.argv[1]).resolve() + if not result_dir.is_dir(): + raise SystemExit(f"result directory not found: {result_dir}") + render(result_dir) + + +if __name__ == "__main__": + main() diff --git a/test/threehost/report-template.html b/test/threehost/report-template.html new file mode 100644 index 0000000..ddb9c41 --- /dev/null +++ b/test/threehost/report-template.html @@ -0,0 +1,1148 @@ + + + + + + + Model‑Velo 三机性能实验报告 + + + +
+ + +
+
+
+
Three‑host benchmark dossier · 2026‑07
+

三机性能实验与热路径证据。

+

客户端 → Model‑Velo → Go 假上游;容量、延迟、依赖池和错误码使用同一 UTC 时间窗。

+
+
+ +
+ 01 / decision line +

当前最可辩护的结论

+
+
+ 0 RPS 严格运行线。 +

严格口径为成功率 ≥ 99.9%、k6 丢弃迭代为 0,并同时观察 P99。

+
+
+
357.56
+
ms · strict point P99
+

低负载 P50 增量 未运行;流式首内容 P50 增量 未运行。

+
+
+
+
0 RPS5007501,0001,5002,000
+
+ + + +
+
+
0–500稳定区
+
500–750严格 SLO 区
+
750–1,000边缘区
+
1,500+过载塌陷区
+
+
+
+
+ +
+
+
02 / load envelope

容量、吞吐与尾延迟

+

闭环并发回答“固定并发下能跑多快”,开环到达率回答“给定业务流量能不能准时接住”。生产容量线应以后者为主,前者用于观察调度和平台形状。

+
+
+
+

闭环网关吞吐

+

c=1…256,三次运行取中位数。吞吐非单调,c=128 后约 1.25k RPS;这不是稳定 SLO。

+
+
+
+

闭环 P99

+

并发上升后排队长尾快速放大:c=256 的 P99 比直连多约 364 ms。

+
+
+
+
+
+

开环:目标、实际与丢弃

+
+
+
+

开环:成功率与 P99

+
+
+
+
+ 1,500 RPS 的问题发生在上游之前,但本轮证据还不足以定位代码。 + 一个代表性 case 中,42,938 个 HTTP 请求只有 24,462 个到达假上游并成功,另有 18,476 个 5xx,P99 约 3.0 秒。假上游最大 active 只有个位数;这排除了“假 LLM 被打满”,却不能在缺少与主运行重叠的网关日志、容器资源和 Prometheus 时序时继续下结论。 +
+
+

直连基线 vs 网关增量

+

直连假上游在高并发时把客户端 CPU 打到 99% 以上;网关 case 的客户端 CPU 峰值仅约 10–12%。因此,摘要里的“客户端可能限制容量”只适用于直连高吞吐基线,不足以解释网关约 1.25k RPS 的平台。

+
并发直连 RPS网关 RPS网关/直连P50 增量P99 增量
+
+
+ +
+
+
03 / server-sent events

SSE 回放与逐 Chunk 观测

+

500 个流式请求、并发 20、三次重复。独立 stream loader 记录 headers、首事件、首内容、总耗时和 chunk 间隔,不只看 k6 的普通 HTTP duration。

+
+
+
+
+

首内容与总完成时间

+
+
+
+

怎么解读

+

首内容 P50 增量只有约 5.37 ms,说明正常路径的转发、SSE 解析与首 chunk 提交成本较小;但首内容 P99 增量约 50 ms,总耗时 P99 也多约 52 ms,尾部抖动值得继续用网关 CPU profile 和 per-stage timing 分解。

+

inter‑chunk P99 只增加约 0.60 ms,说明流一旦开始,持续转发没有明显积累性停顿。三轮 3,000 个直连/网关流请求全部完整成功。

+
边界已经测到。可靠性脚本同时覆盖首 chunk 前非法事件返回 502,以及首 chunk 已提交后断流不再 Fallback;这两条是 SSE 网关最容易写错的状态边界。
+
+
+
+ +
+
+
04 / failure semantics

缓存、故障、队列与 30 分钟耐久

+

这一组不是单纯追求成功率:error10 和 queue overload 会有意制造可接受错误。报告将“测试断言成功”和“业务 HTTP 2xx”分开。

+
+
+
+
+

17 项可靠性断言

+
    +
    +
    +

    三条最重要的诊断

    +

    重试放大:error10 的上游调用约为客户端请求的 1.20×。由于错误由 request ID 确定,重试不会“治愈”这 10%,只会增加上游压力。

    +

    Queue 边界:1,000 RPS 过载 case 中约一半请求成功,上游最大 active 恰为 256;被拒请求没有继续打到上游,边界有效,但 P99 已接近 3 秒。

    +

    耐久错误:500 RPS、30 分钟共 900,001 请求,518 个 5xx,业务成功率 99.942%。它通过 99.9%,但没有达到 99.99%。

    +
    +
    +
    + +
    +
    +
    05 / hot path

    请求阶段、依赖池与错误码

    +

    每个 case 使用开始前和结束后的 Prometheus 累计值做差。阶段分位数来自 Histogram;reliability 包含 Queue、Provider、Retry 和 Fallback,因此不能与其子阶段相加。

    +
    +
    + +
    +
    +
    +
    +

    阶段 P99

    +
    +
    +
    +

    阶段分位数

    +
    +
    阶段样本avg msP50 msP95 msP99 ms
    +
    +
    +
    +
    +

    该 case 的网关错误码

    +
    +
    HTTP内部错误码次数
    +
    +
    +
    +
    + +
    +
    +
    06 / load generator

    客户端资源:峰值属于谁

    +

    客户端有 5,835 个约 1 秒采样。服务器文件也已收到,但它们分别止于主运行前或始于主运行后,必须先做时间窗口审计。

    +
    +
    +
    +
    +

    按阶段 CPU 峰值

    +
    +
    +
    +

    按阶段资源明细

    +
    阶段样本CPU avgCPU max内存 maxk6 RSS max
    +
    +
    +
    +

    服务器采集窗口对齐

    +
    文件类型与主运行关系采样开始 UTC结束 UTC资源峰值
    +
    +
    服务器侧文件存在,但没有覆盖主运行。
    +
    + +
    +
    +
    07 / evidence audit

    哪些结论能说,哪些还不能

    +

    “文件不存在”不是“指标为零”。尤其是 Usage:自动摘要把缺失证据显示成 observed=0,但这不能解释为 243 万事件全部丢失。

    +
    +
    +
    +
    +

    四次运行痕迹

    +
    +
    +
    +

    证据审计结论

    +

    本轮是一次完整、可用于建立容量包络的客户端侧实验,但还不是一次可直接驱动代码改动的服务器侧性能剖析。

    +

    下一轮必须让网关/上游容器 stats、Prometheus、gateway evidence、错误响应分类和网关日志真正覆盖客户端主运行窗口。做到后,1,500 RPS 的 5xx 才能直接映射到 Queue wait timeout、Breaker、Redis、DB pool 或 Usage 链路。

    +

    主运行使用提交 ,且客户端 worktree=dirty。当前工作区之后的参数改动不属于这次实测。

    +
    +
    +
    + +
    +
    +
    08 / external reference

    和其他网关怎么比

    +

    唯一公平的排名来自同一硬件、同一客户端、同一假上游、同一功能开关。本节只提供公开结果的方向性坐标,不宣布跨环境“胜负”。

    +
    +
    +
    +

    独立 CI 的低负载开销参考

    +

    外部项目在同一 neutral CI runner + mock upstream 上测得 Bifrost、Portkey OSS 和 LiteLLM;Model‑Velo 条目来自本次三机 c=1 P50 直连差值,拓扑不同,用斜纹标记。

    +
    +

    方向上,Model‑Velo 3.12 ms 位于该外部 Portkey 2.69 ms 与 LiteLLM 5.41 ms 之间,明显高于 Bifrost 0.56 ms;这只是“值得继续优化”的信号,不是同场排名。

    +
    +
    +

    厂商自报与项目公开状态

    + + + + + + + + + +
    网关公开数字能否直接比
    Bifrostt3.medium 2 vCPU / 4GB,5k RPS、100% 成功;自报内部 overhead 59µs不能 排除项和 2.12s 上游不同
    Portkey OSS项目 README 声称 <1ms;独立 CI 复现为 2.69ms方向参考
    LiteLLM本次检索未找到官方、可复现、同类自托管 RPS 基线;独立 CI 为 5.41ms方向参考
    GoModel官方仓库未发布可复现 benchmark 数字暂无数据
    Model‑Velo本次:严格 750 RPS;低负载 P50 +3.12ms;SSE TTFT P50 +5.37ms本轮实测
    +
    +
    +
    为什么不把 Bifrost 的 5,000 RPS 画进同一吞吐柱图?Bifrost 页面使用 0.13KB 请求、1.37KB 响应、平均端到端 2.12s,并将 JSON marshalling 和 HTTP 调用排除在“59µs overhead”之外;本次 Model‑Velo 是跨三台私网服务器、完整鉴权/限流/缓存/可靠性/Usage 路径。放在一根柱图里会制造虚假的精确性。
    +

    来源与核对日期(2026‑07‑27)

    + +
    + +
    +
    +
    09 / reproducibility

    环境、配置与测量口径

    +

    本报告对应测试提交,不对应当前未提交工作区。所有时间为 UTC;私网地址仅用于复盘三机拓扑。

    +
    +
    +
    +

    测试矩阵

    +
    +
    +
    +

    客户端环境

    +
    +
    +
    +
    查看完整 client-metadata.txt
    +
    + +
    +
    +
    10 / complete case ledger

    全部测试 Case

    +

    包含 smoke、warmup、54 个闭环容量、21 个开环 rate、payload、cache、ramp、burst、fault、queue、endurance、reliability 和 6 个 SSE detail。数值保留到足以回查原始 JSON。

    +
    +
    + + + +
    +
    + + + +
    Case阶段目标负载重复请求RPS成功P50P95P99Dropped上游 calls原始文件
    +
    +

    +
    + +
    +
    +
    11 / raw evidence

    原始产物索引

    +

    主运行目录内的每个文件都列在这里。日志和 JSON 链接使用相对路径,保持本 HTML 与结果目录一起移动即可打开。

    +
    +
    +
    +
    路径大小UTC 修改时间
    +
    +

    +
    查看本报告内嵌数据说明
    summary.html 内嵌 summary.json、client-metadata.txt、client-stats.jsonl 的解析结果、运行尝试目录摘要、可靠性断言名称和产物索引。原始 k6 日志与每个 case 的 JSON 不重复嵌入,使用同目录相对链接访问。
    +
    + +
    +

    Generated · Model‑Velo Performance Lab · 本报告不构成生产 SLA。

    +
    +
    +
    +
    + + + + + diff --git a/test/threehost/run-complete-client.sh b/test/threehost/run-complete-client.sh index 10dacb9..bc66125 100644 --- a/test/threehost/run-complete-client.sh +++ b/test/threehost/run-complete-client.sh @@ -5,7 +5,7 @@ script_dir=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd) repo_root=$(CDPATH= cd -- "$script_dir/../.." && pwd) env_file=${1:-"$script_dir/benchmark.env"} -if [[ ! -f "$env_file" ]]; then +if [[ "$env_file" != "-" && ! -f "$env_file" ]]; then printf 'missing benchmark environment file: %s\n' "$env_file" >&2 exit 1 fi @@ -17,8 +17,10 @@ for command_name in curl git k6 python3; do done set -a -# shellcheck disable=SC1090 -. "$env_file" +if [[ "$env_file" != "-" ]]; then + # shellcheck disable=SC1090 + . "$env_file" +fi set +a : "${GATEWAY_URL:?GATEWAY_URL is required}" @@ -62,6 +64,7 @@ RATE_LIMIT_TEST_RATE=${RATE_LIMIT_TEST_RATE:-100} RATE_LIMIT_TEST_DURATION=${RATE_LIMIT_TEST_DURATION:-2m} RESULTS_ROOT=${RESULTS_ROOT:-"$repo_root/test-results/threehost"} SAVE_RAW_METRICS=${SAVE_RAW_METRICS:-false} +TEST_PROFILE=${TEST_PROFILE:-complete} RUN_CAPACITY=${RUN_CAPACITY:-true} RUN_WARMUP=${RUN_WARMUP:-true} @@ -192,6 +195,7 @@ trap 'cleanup; summarize' EXIT { printf 'run_id=%s\n' "$run_id" + printf 'test_profile=%s\n' "$TEST_PROFILE" printf 'request_prefix=%s\n' "$request_prefix" printf 'commit=%s\n' "$(git -C "$repo_root" rev-parse HEAD)" printf 'started_at=%s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" diff --git a/test/threehost/run-diagnostic-client.sh b/test/threehost/run-diagnostic-client.sh new file mode 100644 index 0000000..07bd578 --- /dev/null +++ b/test/threehost/run-diagnostic-client.sh @@ -0,0 +1,45 @@ +#!/usr/bin/env bash +set -euo pipefail + +script_dir=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd) +env_file=${1:-"$script_dir/benchmark.env"} + +if [[ ! -f "$env_file" ]]; then + printf 'missing benchmark environment file: %s\n' "$env_file" >&2 + exit 1 +fi + +set -a +# shellcheck disable=SC1090 +. "$env_file" +set +a + +export REPETITIONS=1 +export TEST_PROFILE=diagnostic +export WARMUP_RATE=100 +export WARMUP_DURATION=15s +export RATE_SWEEP="500 750 1000 1500" +export RATE_DURATION=2m +export PRE_ALLOCATED_VUS=512 +export MAX_VUS=4096 +export ENDURANCE_RATE=500 +export ENDURANCE_DURATION=10m +export ENDURANCE_PRE_ALLOCATED_VUS=512 +export ENDURANCE_MAX_VUS=2048 + +export RUN_CAPACITY=false +export RUN_WARMUP=true +export RUN_RATE_SWEEP=true +export RUN_STREAM_DETAIL=false +export RUN_PAYLOAD=false +export RUN_CACHE=false +export RUN_RAMP=false +export RUN_BURST=false +export RUN_FAULT=false +export RUN_QUEUE_OVERLOAD=false +export RUN_ENDURANCE=true +export RUN_RELIABILITY=false +export RUN_RATE_LIMIT=false +export SAVE_RAW_METRICS=false + +exec bash "$script_dir/run-complete-client.sh" - diff --git a/test/threehost/summarize.py b/test/threehost/summarize.py index 2aaf143..35542b0 100644 --- a/test/threehost/summarize.py +++ b/test/threehost/summarize.py @@ -7,6 +7,7 @@ import statistics import sys from collections import defaultdict +from datetime import datetime from pathlib import Path @@ -24,6 +25,35 @@ def number(value, default=0.0): return default +def parse_timestamp(value): + if not value: + return None + try: + return datetime.fromisoformat(str(value).replace("Z", "+00:00")) + except ValueError: + return None + + +def case_window(cases): + starts = [ + timestamp + for timestamp in (parse_timestamp(case.get("started_at")) for case in cases) + if timestamp is not None + ] + ends = [ + timestamp + for timestamp in (parse_timestamp(case.get("ended_at")) for case in cases) + if timestamp is not None + ] + return (min(starts) if starts else None, max(ends) if ends else None) + + +def in_window(timestamp, started_at, ended_at): + if timestamp is None or started_at is None or ended_at is None: + return True + return started_at <= timestamp <= ended_at + + def metric_value(metrics, name, field="value"): return number(metrics.get(name, {}).get(field)) @@ -310,18 +340,25 @@ def parse_bytes(value): return amount * base ** powers.get(unit, 0) -def read_resources(result_dir): +def read_resources(result_dir, started_at=None, ended_at=None): samples = defaultdict(list) - for path in result_dir.rglob("*-stats.jsonl"): + for path in result_dir.rglob("*stats*.jsonl"): + if path.name == "client-stats.jsonl": + continue try: lines = path.read_text(encoding="utf-8").splitlines() except OSError: continue for line in lines: try: - stats = json.loads(line).get("stats", {}) + payload = json.loads(line) except json.JSONDecodeError: continue + if not in_window( + parse_timestamp(payload.get("timestamp")), started_at, ended_at + ): + continue + stats = payload.get("stats", {}) name = stats.get("Name") or stats.get("Container") or stats.get("ID") if not name: continue @@ -375,6 +412,7 @@ def read_client_resource(result_dir): return {} cpu = [number(sample.get("cpu_pct")) for sample in samples] + cpu_peak_sample = max(samples, key=lambda sample: number(sample.get("cpu_pct"))) used_memory = [ max( 0, @@ -387,6 +425,7 @@ def read_client_resource(result_dir): "samples": len(samples), "cpu_avg_pct": statistics.fmean(cpu), "cpu_max_pct": max(cpu), + "cpu_peak_at": cpu_peak_sample.get("timestamp", ""), "memory_used_avg_mb": statistics.fmean(used_memory) / (1024 * 1024), "memory_used_max_mb": max(used_memory) / (1024 * 1024), "load_1m_max": max(number(sample.get("load_1m")) for sample in samples), @@ -401,15 +440,26 @@ def read_client_resource(result_dir): PROMETHEUS_LINE = re.compile( - r"^(model_velo_[A-Za-z0-9_:]+(?:\{[^}]*\})?)\s+" + r"^((?:model_velo_|go_|process_)[A-Za-z0-9_:]+(?:\{[^}]*\})?)\s+" r"(-?(?:[0-9]+(?:\.[0-9]*)?|\.[0-9]+)(?:[eE][+-]?[0-9]+)?)$" ) +PROMETHEUS_SERIES = re.compile( + r"^((?:model_velo_|go_|process_)[A-Za-z0-9_:]+)(?:\{(.*)\})?$" +) +PROMETHEUS_LABEL = re.compile(r'([A-Za-z_][A-Za-z0-9_]*)="((?:\\.|[^"])*)"') -def read_prometheus(result_dir): +def read_prometheus(result_dir, started_at=None, ended_at=None): series = {} + points = defaultdict(list) for path in result_dir.rglob("*.promlog"): + snapshot_at = None for line in path.read_text(encoding="utf-8", errors="replace").splitlines(): + if line.startswith("# snapshot "): + snapshot_at = parse_timestamp(line.removeprefix("# snapshot ").strip()) + continue + if not in_window(snapshot_at, started_at, ended_at): + continue match = PROMETHEUS_LINE.match(line.strip()) if not match: continue @@ -417,12 +467,369 @@ def read_prometheus(result_dir): value = number(match.group(2)) current = series.setdefault( key, - {"source": path.name, "series": match.group(1), "samples": 0, "max": value}, + { + "source": path.name, + "series": match.group(1), + "samples": 0, + "first": value, + "min": value, + "max": value, + }, ) current["samples"] += 1 current["last"] = value + current["min"] = min(current["min"], value) current["max"] = max(current["max"], value) - return sorted(series.values(), key=lambda item: (item["source"], item["series"])) + if snapshot_at is not None: + points[key].append((snapshot_at, value)) + output = sorted(series.values(), key=lambda item: (item["source"], item["series"])) + for item in output: + item["delta"] = max(0.0, item.get("last", 0.0) - item["first"]) + return output, points + + +def prometheus_identity(key): + _, separator, series = key.partition(":") + if not separator: + return "", {} + match = PROMETHEUS_SERIES.match(series) + if not match: + return "", {} + labels = {} + for label in PROMETHEUS_LABEL.finditer(match.group(2) or ""): + try: + labels[label.group(1)] = json.loads(f'"{label.group(2)}"') + except json.JSONDecodeError: + labels[label.group(1)] = label.group(2) + return match.group(1), labels + + +def window_delta(samples, started_at, ended_at): + if not samples or started_at is None or ended_at is None: + return 0.0 + before = None + after = None + first_inside = None + last_inside = None + for timestamp, value in samples: + if timestamp <= started_at: + before = value + if started_at <= timestamp <= ended_at: + if first_inside is None: + first_inside = value + last_inside = value + if timestamp >= ended_at: + after = value + break + first = before if before is not None else first_inside + last = after if after is not None else last_inside + if first is None or last is None: + return 0.0 + return max(0.0, last - first) + + +def window_max(samples, started_at, ended_at): + values = [ + value + for timestamp, value in samples + if started_at is not None + and ended_at is not None + and started_at <= timestamp <= ended_at + ] + return max(values) if values else 0.0 + + +def prometheus_counter( + points, metric, started_at, ended_at, labels=None, source=None +): + total = 0.0 + labels = labels or {} + for key, samples in points.items(): + if source and not key.startswith(f"{source}:"): + continue + name, series_labels = prometheus_identity(key) + if name != metric: + continue + if any(series_labels.get(name) != value for name, value in labels.items()): + continue + total += window_delta(samples, started_at, ended_at) + return total + + +def prometheus_gauge_max( + points, metric, started_at, ended_at, labels=None, source=None +): + maximum = 0.0 + labels = labels or {} + for key, samples in points.items(): + if source and not key.startswith(f"{source}:"): + continue + name, series_labels = prometheus_identity(key) + if name != metric: + continue + if any(series_labels.get(name) != value for name, value in labels.items()): + continue + maximum = max(maximum, window_max(samples, started_at, ended_at)) + return maximum + + +def histogram_quantile(buckets, quantile): + if not buckets: + return 0.0 + ordered = sorted(buckets.items(), key=lambda item: item[0]) + count = ordered[-1][1] + if count <= 0: + return 0.0 + target = count * quantile + previous_bound = 0.0 + previous_count = 0.0 + for bound, cumulative in ordered: + if cumulative < target: + if math.isfinite(bound): + previous_bound = bound + previous_count = cumulative + continue + if not math.isfinite(bound): + return previous_bound + bucket_count = cumulative - previous_count + if bucket_count <= 0: + return bound + fraction = (target - previous_count) / bucket_count + return previous_bound + (bound - previous_bound) * fraction + return ordered[-1][0] + + +STAGE_ORDER = { + stage: index + for index, stage in enumerate( + ( + "authentication", + "usage_begin", + "authorization", + "rate_limit", + "route_plan", + "quota_reserve", + "cache_lookup", + "provider_queue", + "provider_call", + "reliability", + "cache_store", + "quota_settle", + "usage_finalize", + ) + ) +} + + +def case_stage_metrics(points, started_at, ended_at): + stages = defaultdict( + lambda: {"buckets": defaultdict(float), "sum_seconds": 0.0, "count": 0.0} + ) + for key, samples in points.items(): + metric, labels = prometheus_identity(key) + stage = labels.get("stage", "") + if not stage: + continue + delta = window_delta(samples, started_at, ended_at) + if metric == "model_velo_request_stage_duration_seconds_bucket": + raw_bound = labels.get("le", "") + bound = math.inf if raw_bound == "+Inf" else number(raw_bound, math.nan) + if not math.isnan(bound): + stages[stage]["buckets"][bound] += delta + elif metric == "model_velo_request_stage_duration_seconds_sum": + stages[stage]["sum_seconds"] += delta + elif metric == "model_velo_request_stage_duration_seconds_count": + stages[stage]["count"] += delta + + output = [] + for stage, values in stages.items(): + count = values["count"] + if count <= 0: + continue + output.append( + { + "stage": stage, + "count": int(count), + "avg_ms": values["sum_seconds"] / count * 1000, + "p50_ms": histogram_quantile(values["buckets"], 0.50) * 1000, + "p95_ms": histogram_quantile(values["buckets"], 0.95) * 1000, + "p99_ms": histogram_quantile(values["buckets"], 0.99) * 1000, + } + ) + return sorted( + output, + key=lambda item: (STAGE_ORDER.get(item["stage"], len(STAGE_ORDER)), item["stage"]), + ) + + +def case_error_counts(points, started_at, ended_at): + counts = defaultdict(float) + for key, samples in points.items(): + metric, labels = prometheus_identity(key) + if metric != "model_velo_http_errors_total": + continue + count = window_delta(samples, started_at, ended_at) + if count > 0: + counts[(labels.get("status", ""), labels.get("code", ""))] += count + return [ + {"status": status, "code": code, "count": int(count)} + for (status, code), count in sorted( + counts.items(), key=lambda item: (-item[1], item[0]) + ) + ] + + +def build_performance_diagnostics(points, rows): + diagnostics = [] + for row in rows: + if row.get("target") != "gateway": + continue + started_at = parse_timestamp(row.get("started_at")) + ended_at = parse_timestamp(row.get("ended_at")) + if started_at is None or ended_at is None: + continue + duration_seconds = max(0.0, (ended_at - started_at).total_seconds()) + stages = case_stage_metrics(points, started_at, ended_at) + diagnostics.append( + { + "case": row.get("case", ""), + "phase": row.get("phase", ""), + "load": row.get("load", ""), + "duration_seconds": duration_seconds, + "stages": stages, + "errors": case_error_counts(points, started_at, ended_at), + "postgres": { + "waits": int( + prometheus_counter( + points, + "model_velo_postgres_waits_total", + started_at, + ended_at, + ) + ), + "wait_ms": prometheus_counter( + points, + "model_velo_postgres_wait_duration_seconds_total", + started_at, + ended_at, + ) + * 1000, + "in_use_max": int( + prometheus_gauge_max( + points, + "model_velo_postgres_connections", + started_at, + ended_at, + {"state": "in_use"}, + ) + ), + "open_max": int( + prometheus_gauge_max( + points, + "model_velo_postgres_connections", + started_at, + ended_at, + {"state": "open"}, + ) + ), + }, + "redis": { + "waits": int( + prometheus_counter( + points, + "model_velo_redis_pool_events_total", + started_at, + ended_at, + {"event": "wait"}, + ) + ), + "timeouts": int( + prometheus_counter( + points, + "model_velo_redis_pool_events_total", + started_at, + ended_at, + {"event": "timeout"}, + ) + ), + "wait_ms": prometheus_counter( + points, + "model_velo_redis_pool_wait_duration_seconds_total", + started_at, + ended_at, + ) + * 1000, + "pending_max": int( + prometheus_gauge_max( + points, + "model_velo_redis_pool_connections", + started_at, + ended_at, + {"state": "pending"}, + ) + ), + "total_max": int( + prometheus_gauge_max( + points, + "model_velo_redis_pool_connections", + started_at, + ended_at, + {"state": "total"}, + ) + ), + }, + "runtime": { + "process_cpu_avg_pct": ( + prometheus_counter( + points, + "process_cpu_seconds_total", + started_at, + ended_at, + source="gateway-metrics.promlog", + ) + / duration_seconds + * 100 + if duration_seconds > 0 + else 0.0 + ), + "rss_max_mb": prometheus_gauge_max( + points, + "process_resident_memory_bytes", + started_at, + ended_at, + source="gateway-metrics.promlog", + ) + / (1024 * 1024), + "heap_max_mb": prometheus_gauge_max( + points, + "go_memstats_heap_alloc_bytes", + started_at, + ended_at, + source="gateway-metrics.promlog", + ) + / (1024 * 1024), + "goroutines_max": int( + prometheus_gauge_max( + points, + "go_goroutines", + started_at, + ended_at, + source="gateway-metrics.promlog", + ) + ), + "gc_cycles": int( + prometheus_counter( + points, + "go_gc_duration_seconds_count", + started_at, + ended_at, + source="gateway-metrics.promlog", + ) + ), + }, + } + ) + return diagnostics def read_usage_evidence(result_dir): @@ -531,10 +938,34 @@ def capacity_findings(capacity, rate, comparisons, stream_pairs): def actionable_findings(resources, client_resource, prometheus, rows): findings = [] if client_resource.get("cpu_max_pct", 0) >= 90: - findings.append( - f"Client host CPU reached {client_resource['cpu_max_pct']:.1f}%; " - "capacity results may be load-generator limited." + peak_at = parse_timestamp(client_resource.get("cpu_peak_at")) + peak_case = next( + ( + row + for row in rows + if peak_at is not None + and ( + parse_timestamp(row.get("started_at")) or peak_at + ) <= peak_at + <= (parse_timestamp(row.get("ended_at")) or peak_at) + ), + {}, ) + if peak_case.get("target") == "direct": + findings.append( + f"Client host CPU reached {client_resource['cpu_max_pct']:.1f}% " + f"during direct baseline {peak_case.get('case')}; " + "that direct-throughput point may be load-generator limited." + ) + else: + findings.append( + f"Client host CPU reached {client_resource['cpu_max_pct']:.1f}%" + + ( + f" during {peak_case.get('case')}." + if peak_case.get("case") + else "." + ) + ) for resource in resources: name = resource["container"].lower() if "gateway" in name and resource["cpu_max_pct"] >= 90: @@ -586,6 +1017,61 @@ def actionable_findings(resources, client_resource, prometheus, rows): return findings +def performance_findings(diagnostics): + findings = [] + diagnosed = [item for item in diagnostics if item["stages"]] + if not diagnosed: + return findings + + pre_provider_names = { + "authentication", + "usage_begin", + "authorization", + "rate_limit", + "route_plan", + "quota_reserve", + } + representative = max( + diagnosed, + key=lambda item: ( + number(item.get("load")), + item.get("duration_seconds", 0), + ), + ) + pre_provider = [ + stage + for stage in representative["stages"] + if stage["stage"] in pre_provider_names + ] + if pre_provider: + slowest = max(pre_provider, key=lambda stage: stage["p99_ms"]) + findings.append( + f"{representative['case']} slowest pre-provider stage was " + f"{slowest['stage']} at P99 {slowest['p99_ms']:.2f} ms." + ) + + postgres_waits = sum(item["postgres"]["waits"] for item in diagnosed) + postgres_wait_ms = sum(item["postgres"]["wait_ms"] for item in diagnosed) + redis_waits = sum(item["redis"]["waits"] for item in diagnosed) + redis_timeouts = sum(item["redis"]["timeouts"] for item in diagnosed) + findings.append( + f"Diagnosed cases recorded {postgres_waits} PostgreSQL pool waits " + f"({postgres_wait_ms:.1f} ms cumulative) and {redis_waits} Redis pool waits " + f"with {redis_timeouts} timeouts." + ) + + errors = defaultdict(int) + for item in diagnosed: + for error in item["errors"]: + errors[error["code"]] += error["count"] + if errors: + code, count = max(errors.items(), key=lambda item: item[1]) + findings.append( + f"Most frequent classified gateway error was {code}: {count} responses." + ) + return findings + + def fmt(value, digits=2): if value is None: return "-" @@ -747,6 +1233,85 @@ def write_markdown(path, summary): ] ) + performance = summary.get("performance_diagnostics", []) + stage_rows = [ + [ + item["case"], + stage["stage"], + stage["count"], + fmt(stage["avg_ms"], 3), + fmt(stage["p50_ms"], 3), + fmt(stage["p95_ms"], 3), + fmt(stage["p99_ms"], 3), + ] + for item in performance + for stage in item["stages"] + ] + if stage_rows: + lines.extend( + [ + "## Hot-path stage timing", + "", + "The `reliability` row contains queue, provider calls, retries, and " + "fallbacks; it overlaps the `provider_queue` and `provider_call` rows.", + "", + markdown_table( + ["case", "stage", "count", "avg ms", "P50 ms", "P95 ms", "P99 ms"], + stage_rows, + ), + "", + "### Dependency pools and Go runtime", + "", + markdown_table( + [ + "case", + "PG waits", + "PG wait ms", + "PG in-use max", + "Redis waits", + "Redis wait ms", + "Redis pending max", + "process CPU avg", + "RSS max MB", + "goroutines max", + ], + [ + [ + item["case"], + item["postgres"]["waits"], + fmt(item["postgres"]["wait_ms"], 2), + item["postgres"]["in_use_max"], + item["redis"]["waits"], + fmt(item["redis"]["wait_ms"], 2), + item["redis"]["pending_max"], + fmt(item["runtime"]["process_cpu_avg_pct"], 1) + "%", + fmt(item["runtime"]["rss_max_mb"], 1), + item["runtime"]["goroutines_max"], + ] + for item in performance + ], + ), + "", + ] + ) + errors = [ + [item["case"], error["status"], error["code"], error["count"]] + for item in performance + for error in item["errors"] + ] + if errors: + lines.extend( + [ + "### Gateway error codes", + "", + markdown_table( + ["case", "HTTP status", "error code", "count"], + errors, + ), + "", + ] + ) + resources = summary["resources"] if resources: lines.extend( @@ -871,19 +1436,46 @@ def main(): "reliability", } ] - resources = read_resources(result_dir) + started_at, ended_at = case_window(cases) + resources = read_resources(result_dir, started_at, ended_at) client_resource = read_client_resource(result_dir) - prometheus = read_prometheus(result_dir) + prometheus, prometheus_points = read_prometheus( + result_dir, started_at, ended_at + ) + performance_diagnostics = build_performance_diagnostics( + prometheus_points, k6_rows + ) usage = read_usage_evidence(result_dir) usage_reconciliation = reconcile_usage(k6_rows, stream_rows, usage) warnings = [] - if not resources: + stats_paths = [ + path + for path in result_dir.rglob("*stats*.jsonl") + if path.name != "client-stats.jsonl" + ] + prometheus_paths = list(result_dir.rglob("*.promlog")) + if stats_paths and not resources: + warnings.append( + "Docker stats files exist, but their timestamps do not overlap " + "the benchmark case window." + ) + elif not resources: warnings.append("Missing gateway/upstream Docker stats JSONL files.") if not client_resource: warnings.append("Missing client host resource samples.") - if not prometheus: + if prometheus_paths and not prometheus: + warnings.append( + "Prometheus capture exists, but its timestamps do not overlap " + "the benchmark case window." + ) + elif not prometheus: warnings.append("Missing Prometheus time-series capture.") + elif not any(item["stages"] for item in performance_diagnostics): + warnings.append( + "Prometheus capture does not contain per-stage request histograms; " + "confirm the gateway image includes diagnostic metrics." + ) if not usage: warnings.append("Missing post-run Usage/PostgreSQL/Redis evidence.") elif usage.get("drain", {}).get("state") not in {"", "complete", None}: @@ -918,6 +1510,7 @@ def main(): k6_rows + stream_rows, ) ) + findings.extend(performance_findings(performance_diagnostics)) if not findings: if k6_rows or stream_rows: findings.append( @@ -941,6 +1534,7 @@ def main(): "resources": resources, "client_resource": client_resource, "prometheus": prometheus, + "performance_diagnostics": performance_diagnostics, "usage_evidence": usage, "usage_reconciliation": usage_reconciliation, "findings": findings, @@ -951,7 +1545,13 @@ def main(): encoding="utf-8", ) write_markdown(result_dir / "summary.md", summary) - print(f"wrote {result_dir / 'summary.json'} and {result_dir / 'summary.md'}") + from render_html import render + + render(result_dir) + print( + f"wrote {result_dir / 'summary.json'}, " + f"{result_dir / 'summary.md'}, and {result_dir / 'summary.html'}" + ) if __name__ == "__main__": From ad7dbf834f47dc086267dca0e2671ae9dfbf5db0 Mon Sep 17 00:00:00 2001 From: tonghzhang <949674415@qq.com> Date: Tue, 28 Jul 2026 16:04:07 +0800 Subject: [PATCH 2/9] fix(usage): stabilize and instrument benchmark delivery --- cmd/model-velo-usage-worker/main.go | 7 +- cmd/model-velo/main.go | 15 +- internal/httpapi/auth.go | 20 ++- internal/httpapi/chat.go | 99 +++++++++++- internal/httpapi/chat_test.go | 10 +- internal/httpapi/router_options.go | 5 + internal/httpapi/usage.go | 38 ++++- internal/observability/dependencies.go | 145 +++++++++++++++++ internal/observability/metrics.go | 46 ++++++ internal/observability/observability_test.go | 17 ++ internal/reliability/queue_observer.go | 36 +++++ internal/reliability/tracing.go | 8 + internal/usage/collector.go | 2 + internal/usage/emitter.go | 2 + internal/usage/outbox.go | 158 ++++++++++++------- internal/usage/store.go | 49 +++++- internal/usage/usage_test.go | 41 ++++- internal/usage/worker.go | 72 +++++---- test/threehost/collect-run-evidence.sh | 10 +- 19 files changed, 652 insertions(+), 128 deletions(-) create mode 100644 internal/observability/dependencies.go create mode 100644 internal/reliability/queue_observer.go diff --git a/cmd/model-velo-usage-worker/main.go b/cmd/model-velo-usage-worker/main.go index 81d0653..b9cbe00 100644 --- a/cmd/model-velo-usage-worker/main.go +++ b/cmd/model-velo-usage-worker/main.go @@ -96,7 +96,12 @@ func run() error { if err != nil { return fmt.Errorf("configure usage outbox emitter: %w", err) } - relay, err := usage.NewOutboxRelay(database.ORM(), redisEmitter, usageConfig.WorkerTimeout) + relay, err := usage.NewOutboxRelay( + database.ORM(), + redisEmitter, + usageConfig.Group, + usageConfig.BatchSize, + ) if err != nil { return fmt.Errorf("configure usage outbox relay: %w", err) } diff --git a/cmd/model-velo/main.go b/cmd/model-velo/main.go index ae7199d..714e632 100644 --- a/cmd/model-velo/main.go +++ b/cmd/model-velo/main.go @@ -110,17 +110,8 @@ func run() error { // 装配基础设施、启动 HTTP 服务并等待关闭。 if err != nil { // 响应缓存配置无效。 return fmt.Errorf("configure response cache: %w", err) // 返回缓存组件配置错误。 } - redisUsageEmitter, err := usage.NewRedisEmitter( - redisClient.Native(), - startup.usage.StreamKey, - startup.usage.EmitTimeout, - ) - if err != nil { - return fmt.Errorf("configure usage emitter: %w", err) - } usageEmitter, err := usage.NewDurableEmitter( database.ORM(), - redisUsageEmitter, startup.usage.EmitTimeout, ) if err != nil { @@ -183,6 +174,12 @@ func run() error { // 装配基础设施、启动 HTTP 服务并等待关闭。 if err := metrics.RegisterRuntime(runtimeManager); err != nil { return fmt.Errorf("register runtime metrics: %w", err) } + if err := metrics.RegisterDependencies( + database.SQL(), + redisClient.Native(), + ); err != nil { + return fmt.Errorf("register dependency metrics: %w", err) + } readiness := health.NewChecker( database.SQL(), redisClient.Native(), diff --git a/internal/httpapi/auth.go b/internal/httpapi/auth.go index de24926..40641fc 100644 --- a/internal/httpapi/auth.go +++ b/internal/httpapi/auth.go @@ -5,6 +5,7 @@ import ( "errors" // 判断具体认证错误类型。 "net/http" // 读取请求头并使用 HTTP 状态码。 "strings" // 解析 Authorization 请求头。 + "time" "github.com/gin-gonic/gin" // Gin 中间件与请求上下文。 @@ -44,8 +45,25 @@ func authenticationMiddleware(access AccessController) gin.HandlerFunc { // 创 return // 不再进入后续路由。 } + startedAt := time.Now() identity, authenticateErr := access.Authenticate(c.Request.Context(), plaintext) // 验证 API Key 并取得租户身份。 - if authenticateErr != nil { // API Key 验证失败。 + result := "accepted" + switch { + case c.Request.Context().Err() != nil: + result = "canceled" + case errors.Is(authenticateErr, apikey.ErrInvalidCredential), + errors.Is(authenticateErr, apikey.ErrKeyInactive), + errors.Is(authenticateErr, apikey.ErrKeyRevoked), + errors.Is(authenticateErr, apikey.ErrKeyExpired), + errors.Is(authenticateErr, apikey.ErrTenantInactive): + result = "rejected" + case authenticateErr != nil: + result = "error" + } + routerMetrics(c).RequestStage( + "authentication", result, "", time.Since(startedAt), + ) + if authenticateErr != nil { // API Key 验证失败。 if c.Request.Context().Err() != nil { // 请求已经取消时不再写响应。 return } diff --git a/internal/httpapi/chat.go b/internal/httpapi/chat.go index f0835d3..f7b476f 100644 --- a/internal/httpapi/chat.go +++ b/internal/httpapi/chat.go @@ -163,15 +163,29 @@ func (h chatHandler) complete(c *gin.Context) { return } - if err := h.access.AuthorizeModel( // 检查该租户是否允许使用当前模型。 + startedAt := time.Now() + authorizationErr := h.access.AuthorizeModel( // 检查该租户是否允许使用当前模型。 c.Request.Context(), identity.TenantID, model, - ); err != nil { + ) + authorizationResult := "allowed" + switch { + case c.Request.Context().Err() != nil: + authorizationResult = "canceled" + case errors.Is(authorizationErr, apikey.ErrModelNotAllowed): + authorizationResult = "denied" + case authorizationErr != nil: + authorizationResult = "error" + } + h.metrics.RequestStage( + "authorization", authorizationResult, "", time.Since(startedAt), + ) + if authorizationErr != nil { if c.Request.Context().Err() != nil { // 请求已经取消时不再写响应。 return } - if errors.Is(err, apikey.ErrModelNotAllowed) { // API Key 没有模型权限。 + if errors.Is(authorizationErr, apikey.ErrModelNotAllowed) { // API Key 没有模型权限。 writeAPIError( c, http.StatusForbidden, // 返回 403。 @@ -194,7 +208,22 @@ func (h chatHandler) complete(c *gin.Context) { return } + startedAt = time.Now() limitDecision, err := h.limiter.Allow(c.Request.Context(), identity.TenantID, model) // 查询租户和模型的限流额度。 + limitResult := "allowed" + switch { + case c.Request.Context().Err() != nil: + limitResult = "canceled" + case err != nil: + limitResult = "error" + case !limitDecision.Allowed: + limitResult = "rejected" + case limitDecision.Bypassed: + limitResult = "bypassed" + } + h.metrics.RequestStage( + "rate_limit", limitResult, "", time.Since(startedAt), + ) if err != nil { if c.Request.Context().Err() != nil { // 请求已经取消。 return @@ -257,7 +286,15 @@ func (h chatHandler) complete(c *gin.Context) { c.Request.Context(), active.CacheNamespace, )) + startedAt = time.Now() routePlan, err := active.Routes.Plan(model, requiredCapabilities) // 生成符合模型和能力要求的候选 Provider。 + routeResult := "planned" + if err != nil { + routeResult = "error" + } + h.metrics.RequestStage( + "route_plan", routeResult, "", time.Since(startedAt), + ) if err != nil { if errors.Is(err, routing.ErrCapabilityUnavailable) { // 没有 Provider 支持请求能力。 writeAPIError( @@ -306,7 +343,15 @@ func (h chatHandler) complete(c *gin.Context) { cacheResult := responsecache.Result{Status: responsecache.StatusBypass} // 默认标记为绕过缓存。 if !requestBypassesResponseCache(c.Request) { // 请求头没有 no-store 时才查询缓存。 + startedAt = time.Now() cacheResult, err = h.cache.Lookup(c.Request.Context(), identity.TenantID, model, requestBody) // 用租户、模型和请求体查缓存。 + cacheLookupResult := string(cacheResult.Status) + if err != nil { + cacheLookupResult = "error" + } + h.metrics.RequestStage( + "cache_lookup", cacheLookupResult, "", time.Since(startedAt), + ) if err != nil { if c.Request.Context().Err() != nil { // 请求取消时停止处理。 return @@ -349,11 +394,19 @@ func (h chatHandler) complete(c *gin.Context) { h.metrics.Cache("lookup", string(cacheResult.Status)) usageSession.setCacheStatus(string(cacheResult.Status)) + startedAt = time.Now() execution, failure := active.Chat.Execute(c.Request.Context(), reliability.ExecutionInput{ // 执行 Provider 调用和回退。 RequestID: requestIDFromContext(c.Request.Context()), // 当前网关请求 ID。 Request: request, // 解析后的聊天请求。 Plan: routePlan, // 路由候选计划。 }) + executionResult := "success" + if failure != nil { + executionResult = string(failure.Category) + } + h.metrics.RequestStage( + "reliability", executionResult, "", time.Since(startedAt), + ) if failure != nil { h.observeFailure(requestIDFromContext(c.Request.Context()), failure) usageSession.recordFailure(failure) @@ -379,13 +432,22 @@ func (h chatHandler) complete(c *gin.Context) { usageSession.recordExecution(execution) usageSession.observe(responseBody) if cacheResult.Status == responsecache.StatusMiss && execution.Fallbacks == 0 { // 只有缓存未命中且未切换 Provider 时才缓存。 - if err := h.cache.Store( + startedAt = time.Now() + err := h.cache.Store( c.Request.Context(), identity.TenantID, model, requestBody, responseBody, - ); err != nil { + ) + cacheStoreResult := "stored" + if err != nil { + cacheStoreResult = "error" + } + h.metrics.RequestStage( + "cache_store", cacheStoreResult, "", time.Since(startedAt), + ) + if err != nil { if c.Request.Context().Err() != nil { // 请求取消时停止。 return } @@ -442,6 +504,7 @@ func (h chatHandler) reserveQuota( if request.N != nil && *request.N > 1 { output *= int64(*request.N) } + startedAt := time.Now() decision, err := h.quota.Reserve(c.Request.Context(), quota.ReserveInput{ GroupID: session.eventID, TenantID: tenantID, @@ -450,6 +513,24 @@ func (h chatHandler) reserveQuota( EstimatedOutputTokens: output, Plan: plan, }) + result := "allowed" + switch { + case c.Request.Context().Err() != nil: + result = "canceled" + case errors.Is(err, quota.ErrExceeded): + result = "denied" + case err != nil: + result = "error" + case decision.Exceeded: + result = "allowed_overage" + case len(decision.Alerts) > 0: + result = "allowed_alert" + case decision.AppliedPolicies == 0: + result = "no_policy" + } + h.metrics.RequestStage( + "quota_reserve", result, "", time.Since(startedAt), + ) if err != nil { if errors.Is(err, quota.ErrExceeded) { h.metrics.Quota("denied") @@ -598,11 +679,19 @@ func (h chatHandler) stream( return } + startedAt := time.Now() prepared, failure := orchestrator.OpenStream(c.Request.Context(), reliability.ExecutionInput{ // 打开上游流并取得首个事件。 RequestID: requestIDFromContext(c.Request.Context()), Request: request, Plan: plan, }) + result := "success" + if failure != nil { + result = string(failure.Category) + } + h.metrics.RequestStage( + "reliability", result, "", time.Since(startedAt), + ) if failure != nil { h.observeFailure(requestIDFromContext(c.Request.Context()), failure) usageSession.recordFailure(failure) diff --git a/internal/httpapi/chat_test.go b/internal/httpapi/chat_test.go index 698d8bf..248ccba 100644 --- a/internal/httpapi/chat_test.go +++ b/internal/httpapi/chat_test.go @@ -158,8 +158,11 @@ func TestChatCompletionUsageLifecycle(t *testing.T) { } event := singleUsageEvent(t, emitter) + if err := event.Validate(); err != nil { + t.Fatalf("cache usage event validation error = %v", err) + } if event.Status != usage.StatusCacheHit || - event.CacheStatus != string(responsecache.StatusHit) || + event.CacheStatus != "hit" || event.UsageSource != usage.UsageSourceCacheReplay || event.ProviderID != "" || event.Attempts != 0 || @@ -236,11 +239,14 @@ func TestChatCompletionUsageLifecycle(t *testing.T) { } event := singleUsageEvent(t, emitter) + if err := event.Validate(); err != nil { + t.Fatalf("stream usage event validation error = %v", err) + } if event.Status != usage.StatusStreamCompleted || !event.Stream || event.FirstTokenMS == nil || event.UsageSource != usage.UsageSourceProvider || - event.CacheStatus != string(responsecache.StatusBypass) || + event.CacheStatus != "bypass" || event.Usage == nil || event.Usage.Total != 7 { t.Fatalf("stream usage event = %#v", event) diff --git a/internal/httpapi/router_options.go b/internal/httpapi/router_options.go index aeccafc..13eb3be 100644 --- a/internal/httpapi/router_options.go +++ b/internal/httpapi/router_options.go @@ -20,6 +20,7 @@ import ( "model-velo/internal/gateway" "model-velo/internal/observability" "model-velo/internal/quota" + "model-velo/internal/reliability" ) type ReadinessChecker interface { @@ -104,6 +105,9 @@ func requestSummaryMiddleware( c.Header("request-id", requestIDFromContext(c.Request.Context())) } c.Set("model-velo.metrics", metrics) + c.Request = c.Request.WithContext( + reliability.WithQueueObserver(c.Request.Context(), metrics), + ) finishMetrics := metrics.BeginRequest() c.Next() @@ -116,6 +120,7 @@ func requestSummaryMiddleware( status := c.Writer.Status() duration := time.Since(startedAt) finishMetrics(route, c.Request.Method, status, isStream, duration) + metrics.HTTPError(route, status, apiErrorCode(c)) attributes := []any{ "request_id", requestIDFromContext(c.Request.Context()), diff --git a/internal/httpapi/usage.go b/internal/httpapi/usage.go index 38672a6..e73cc55 100644 --- a/internal/httpapi/usage.go +++ b/internal/httpapi/usage.go @@ -36,6 +36,14 @@ func newUsageSession( stream bool, metrics *observability.Metrics, ) (*usageSession, error) { + startedAt := time.Now() + stageResult := "error" + defer func() { + metrics.RequestStage( + "usage_begin", stageResult, "", time.Since(startedAt), + ) + }() + collector, err := usage.NewCollector(usage.NewEventInput{ RequestID: requestIDFromContext(ctx), TenantID: tenantID, @@ -57,6 +65,7 @@ func newUsageSession( } metrics.UsageDelivery("lifecycle_begin", "success") } + stageResult = "success" return &usageSession{ collector: collector, emitter: emitter, eventID: pending.EventID, metrics: metrics, @@ -94,9 +103,13 @@ func (session *usageSession) finish(c *gin.Context) { settleContext, cancel := context.WithTimeout( context.WithoutCancel(c.Request.Context()), 2*time.Second, ) - if err := session.quota.Settle( + startedAt := time.Now() + err := session.quota.Settle( settleContext, session.reservationID, session.event, - ); err != nil { + ) + result := "success" + if err != nil { + result = "error" slog.Error( "quota settlement failed", "request_id", session.event.RequestID, @@ -104,19 +117,34 @@ func (session *usageSession) finish(c *gin.Context) { "error", err, ) } + session.metrics.RequestStage( + "quota_settle", result, "", time.Since(startedAt), + ) cancel() } - if _, err := session.emitter.Emit(c.Request.Context(), session.event); err != nil { + startedAt := time.Now() + entryID, err := session.emitter.Emit(c.Request.Context(), session.event) + result := "queued" + if entryID != "" { + result = "published" + } + if err != nil { + result = "deferred" + } + session.metrics.RequestStage( + "usage_finalize", result, "", time.Since(startedAt), + ) + if err != nil { session.metrics.UsageDelivery("finalize", "deferred") slog.Warn( - "usage immediate publish deferred", + "usage finalization deferred", "request_id", session.event.RequestID, "event_id", session.event.EventID, "error", err, ) return } - session.metrics.UsageDelivery("finalize", "published") + session.metrics.UsageDelivery("finalize", result) } func (session *usageSession) attachQuota( diff --git a/internal/observability/dependencies.go b/internal/observability/dependencies.go new file mode 100644 index 0000000..30f1a43 --- /dev/null +++ b/internal/observability/dependencies.go @@ -0,0 +1,145 @@ +package observability + +import ( + "database/sql" + + "github.com/prometheus/client_golang/prometheus" + goredis "github.com/redis/go-redis/v9" +) + +type dependencyCollector struct { + database *sql.DB + redis *goredis.Client + + postgresConnections *prometheus.Desc + postgresWaits *prometheus.Desc + postgresWaitTime *prometheus.Desc + redisConnections *prometheus.Desc + redisPoolEvents *prometheus.Desc + redisWaitTime *prometheus.Desc +} + +func (metrics *Metrics) RegisterDependencies( + database *sql.DB, + redis *goredis.Client, +) error { + if metrics == nil || (database == nil && redis == nil) { + return nil + } + return metrics.registry.Register(newDependencyCollector(database, redis)) +} + +func newDependencyCollector( + database *sql.DB, + redis *goredis.Client, +) *dependencyCollector { + return &dependencyCollector{ + database: database, + redis: redis, + postgresConnections: prometheus.NewDesc( + "model_velo_postgres_connections", + "Current PostgreSQL connection-pool state.", + []string{"state"}, nil, + ), + postgresWaits: prometheus.NewDesc( + "model_velo_postgres_waits_total", + "PostgreSQL connection-pool waits.", + nil, nil, + ), + postgresWaitTime: prometheus.NewDesc( + "model_velo_postgres_wait_duration_seconds_total", + "Total time waiting for a PostgreSQL connection.", + nil, nil, + ), + redisConnections: prometheus.NewDesc( + "model_velo_redis_pool_connections", + "Current Redis connection-pool state.", + []string{"state"}, nil, + ), + redisPoolEvents: prometheus.NewDesc( + "model_velo_redis_pool_events_total", + "Redis connection-pool cumulative outcomes.", + []string{"event"}, nil, + ), + redisWaitTime: prometheus.NewDesc( + "model_velo_redis_pool_wait_duration_seconds_total", + "Total time waiting for a Redis connection.", + nil, nil, + ), + } +} + +func (collector *dependencyCollector) Describe(output chan<- *prometheus.Desc) { + output <- collector.postgresConnections + output <- collector.postgresWaits + output <- collector.postgresWaitTime + output <- collector.redisConnections + output <- collector.redisPoolEvents + output <- collector.redisWaitTime +} + +func (collector *dependencyCollector) Collect(output chan<- prometheus.Metric) { + if collector.database != nil { + stats := collector.database.Stats() + for state, value := range map[string]int{ + "open": stats.OpenConnections, + "in_use": stats.InUse, + "idle": stats.Idle, + "max_open": stats.MaxOpenConnections, + } { + output <- prometheus.MustNewConstMetric( + collector.postgresConnections, + prometheus.GaugeValue, + float64(value), + state, + ) + } + output <- prometheus.MustNewConstMetric( + collector.postgresWaits, + prometheus.CounterValue, + float64(stats.WaitCount), + ) + output <- prometheus.MustNewConstMetric( + collector.postgresWaitTime, + prometheus.CounterValue, + stats.WaitDuration.Seconds(), + ) + } + + if collector.redis == nil { + return + } + stats := collector.redis.PoolStats() + for state, value := range map[string]uint32{ + "total": stats.TotalConns, + "idle": stats.IdleConns, + "pending": stats.PendingRequests, + } { + output <- prometheus.MustNewConstMetric( + collector.redisConnections, + prometheus.GaugeValue, + float64(value), + state, + ) + } + for event, value := range map[string]uint32{ + "hit": stats.Hits, + "miss": stats.Misses, + "timeout": stats.Timeouts, + "wait": stats.WaitCount, + "unusable": stats.Unusable, + "stale": stats.StaleConns, + } { + output <- prometheus.MustNewConstMetric( + collector.redisPoolEvents, + prometheus.CounterValue, + float64(value), + event, + ) + } + output <- prometheus.MustNewConstMetric( + collector.redisWaitTime, + prometheus.CounterValue, + float64(stats.WaitDurationNs)/1e9, + ) +} diff --git a/internal/observability/metrics.go b/internal/observability/metrics.go index 847f7a6..c20fa8b 100644 --- a/internal/observability/metrics.go +++ b/internal/observability/metrics.go @@ -19,6 +19,8 @@ type Metrics struct { registry *prometheus.Registry requests *prometheus.CounterVec requestDuration *prometheus.HistogramVec + requestErrors *prometheus.CounterVec + stageDuration *prometheus.HistogramVec inFlight prometheus.Gauge providerAttempts *prometheus.CounterVec providerDuration *prometheus.HistogramVec @@ -43,6 +45,18 @@ func NewMetrics() *Metrics { Help: "End-to-end HTTP request duration.", Buckets: []float64{0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 120, 300}, }, []string{"route", "method", "status", "stream"}), + requestErrors: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "model_velo_http_errors_total", + Help: "Completed HTTP errors by stable gateway error code.", + }, []string{"route", "status", "code"}), + stageDuration: prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Name: "model_velo_request_stage_duration_seconds", + Help: "Duration of bounded gateway request stages.", + Buckets: []float64{ + 0.0001, 0.00025, 0.0005, 0.001, 0.0025, 0.005, + 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, + }, + }, []string{"stage", "result", "provider"}), inFlight: prometheus.NewGauge(prometheus.GaugeOpts{ Name: "model_velo_http_in_flight", Help: "Current in-flight HTTP requests.", @@ -88,6 +102,8 @@ func NewMetrics() *Metrics { metrics.registry.MustRegister( metrics.requests, metrics.requestDuration, + metrics.requestErrors, + metrics.stageDuration, metrics.inFlight, metrics.providerAttempts, metrics.providerDuration, @@ -98,6 +114,8 @@ func NewMetrics() *Metrics { metrics.auth, metrics.usageDelivery, metrics.quota, + prometheus.NewGoCollector(), + prometheus.NewProcessCollector(prometheus.ProcessCollectorOpts{}), ) return metrics } @@ -163,6 +181,33 @@ func (metrics *Metrics) Authentication(result string) { } } +func (metrics *Metrics) HTTPError(route string, status int, code string) { + if metrics == nil || status < http.StatusBadRequest { + return + } + if code == "" { + code = "unclassified" + } + metrics.requestErrors.WithLabelValues(route, strconv.Itoa(status), code).Inc() +} + +func (metrics *Metrics) RequestStage( + stage, result, provider string, + duration time.Duration, +) { + if metrics != nil { + metrics.stageDuration.WithLabelValues(stage, result, provider). + Observe(duration.Seconds()) + } +} + +func (metrics *Metrics) ObserveQueueWait( + provider, result string, + duration time.Duration, +) { + metrics.RequestStage("provider_queue", result, provider, duration) +} + func (metrics *Metrics) RateLimit(result string) { if metrics != nil { metrics.rateLimit.WithLabelValues(result).Inc() @@ -185,6 +230,7 @@ func (metrics *Metrics) ProviderAttempt( } metrics.providerAttempts.WithLabelValues(provider, result, category).Inc() metrics.providerDuration.WithLabelValues(provider, result).Observe(duration.Seconds()) + metrics.RequestStage("provider_call", result, provider, duration) if retry { metrics.retries.WithLabelValues(provider, category).Inc() } diff --git a/internal/observability/observability_test.go b/internal/observability/observability_test.go index 147bdf3..ae5f96c 100644 --- a/internal/observability/observability_test.go +++ b/internal/observability/observability_test.go @@ -2,12 +2,15 @@ package observability import ( "context" + "database/sql" "net/http" "net/http/httptest" "strings" "testing" "time" + goredis "github.com/redis/go-redis/v9" + "model-velo/internal/usage" ) @@ -16,8 +19,17 @@ func TestMetricsAreBoundedAndProtected(t *testing.T) { if err := metrics.RegisterUsageWorker(fakeUsageWorker{}); err != nil { t.Fatal(err) } + redisClient := goredis.NewClient(&goredis.Options{Addr: "127.0.0.1:1"}) + t.Cleanup(func() { + _ = redisClient.Close() + }) + if err := metrics.RegisterDependencies(&sql.DB{}, redisClient); err != nil { + t.Fatal(err) + } finish := metrics.BeginRequest() finish("/v1/test", http.MethodPost, http.StatusOK, false, 20*time.Millisecond) + metrics.HTTPError("/v1/test", http.StatusServiceUnavailable, "gateway_overloaded") + metrics.RequestStage("authentication", "accepted", "", time.Millisecond) metrics.Authentication("accepted") metrics.RateLimit("allowed") metrics.Cache("lookup", "hit") @@ -53,11 +65,16 @@ func TestMetricsAreBoundedAndProtected(t *testing.T) { payload := response.Body.String() for _, metricName := range []string{ "model_velo_http_requests_total", + "model_velo_http_errors_total", + "model_velo_request_stage_duration_seconds", "model_velo_provider_attempts_total", "model_velo_authentication_total", "model_velo_usage_delivery_total", "model_velo_usage_worker_pending", "model_velo_quota_decisions_total", + "model_velo_postgres_connections", + "model_velo_redis_pool_connections", + "go_goroutines", } { if !strings.Contains(payload, metricName) { t.Errorf("scrape does not contain %s", metricName) diff --git a/internal/reliability/queue_observer.go b/internal/reliability/queue_observer.go new file mode 100644 index 0000000..75eec4c --- /dev/null +++ b/internal/reliability/queue_observer.go @@ -0,0 +1,36 @@ +package reliability + +import ( + "context" + "time" +) + +type QueueObserver interface { + ObserveQueueWait(provider, result string, duration time.Duration) +} + +type queueObserverContextKey struct{} + +func WithQueueObserver(ctx context.Context, observer QueueObserver) context.Context { + if ctx == nil { + ctx = context.Background() + } + if observer == nil { + return ctx + } + return context.WithValue(ctx, queueObserverContextKey{}, observer) +} + +func observeQueueWait( + ctx context.Context, + provider, result string, + duration time.Duration, +) { + if ctx == nil { + return + } + observer, _ := ctx.Value(queueObserverContextKey{}).(QueueObserver) + if observer != nil { + observer.ObserveQueueWait(provider, result, duration) + } +} diff --git a/internal/reliability/tracing.go b/internal/reliability/tracing.go index a944885..66d1e86 100644 --- a/internal/reliability/tracing.go +++ b/internal/reliability/tracing.go @@ -87,17 +87,25 @@ func traceQueueAcquire( attribute.String("gateway.provider.id", providerID), ), ) + startedAt := time.Now() lease, failure := queues.Acquire(ctx, providerID) + duration := time.Since(startedAt) + result := "acquired" if failure == nil { span.SetAttributes(attribute.String("gateway.queue.result", "acquired")) span.SetStatus(codes.Ok, "") } else { + result = string(failure.Queue) + if result == "" { + result = string(failure.Category) + } span.SetAttributes( attribute.String("gateway.queue.result", string(failure.Queue)), attribute.String("gateway.failure.category", string(failure.Category)), ) span.SetStatus(codes.Error, string(failure.Category)) } + observeQueueWait(ctx, providerID, result, duration) span.End() return lease, failure } diff --git a/internal/usage/collector.go b/internal/usage/collector.go index 1158f97..8bb283d 100644 --- a/internal/usage/collector.go +++ b/internal/usage/collector.go @@ -1,6 +1,7 @@ package usage import ( + "strings" "sync" "time" ) @@ -40,6 +41,7 @@ func (collector *Collector) SetCacheStatus(status string) { if collector == nil { return } + status = strings.ToLower(strings.TrimSpace(status)) collector.mu.Lock() defer collector.mu.Unlock() if !collector.finalized && status != "" { diff --git a/internal/usage/emitter.go b/internal/usage/emitter.go index 1824276..fffa25c 100644 --- a/internal/usage/emitter.go +++ b/internal/usage/emitter.go @@ -9,6 +9,8 @@ import ( ) type Emitter interface { + // Emit returns a Redis entry ID for immediate delivery. A durable emitter + // may return an empty ID after the event is safely queued in its outbox. Emit(ctx context.Context, event Event) (string, error) } diff --git a/internal/usage/outbox.go b/internal/usage/outbox.go index 73f07f8..3f3e30d 100644 --- a/internal/usage/outbox.go +++ b/internal/usage/outbox.go @@ -12,10 +12,7 @@ import ( "model-velo/internal/postgres" ) -const ( - defaultOutboxBatchSize = 100 - defaultOutboxRepublishPeriod = 30 * time.Second -) +const defaultOutboxRepublishPeriod = 30 * time.Second // PendingEvent is the durable minimum recorded before an authenticated request // reaches a provider. It intentionally excludes prompts and provider secrets. @@ -36,29 +33,24 @@ type LifecycleEmitter interface { Begin(context.Context, PendingEvent) error } -// DurableEmitter stores request lifecycle state in PostgreSQL before making a -// best-effort immediate delivery to Redis. The worker relays anything left in -// the outbox, so Redis downtime cannot silently discard a finalized event. +// DurableEmitter stores request lifecycle state in PostgreSQL. The worker owns +// Redis publication so online requests do not compete with replay traffic. type DurableEmitter struct { database *gorm.DB - redis *RedisEmitter timeout time.Duration } func NewDurableEmitter( database *gorm.DB, - redis *RedisEmitter, timeout time.Duration, ) (*DurableEmitter, error) { switch { case database == nil: return nil, errors.New("durable usage emitter requires PostgreSQL") - case redis == nil: - return nil, errors.New("durable usage emitter requires Redis") case timeout <= 0: return nil, errors.New("durable usage emitter timeout must be positive") default: - return &DurableEmitter{database: database, redis: redis, timeout: timeout}, nil + return &DurableEmitter{database: database, timeout: timeout}, nil } } @@ -79,6 +71,7 @@ func (emitter *DurableEmitter) Begin(ctx context.Context, pending PendingEvent) StartedAt: pending.StartedAt.UTC(), } result := emitter.database.WithContext(writeContext). + Session(&gorm.Session{SkipDefaultTransaction: true}). Clauses(clause.OnConflict{Columns: []clause.Column{{Name: "event_id"}}, DoNothing: true}). Create(&record) if result.Error != nil { @@ -95,6 +88,7 @@ func (emitter *DurableEmitter) Emit(ctx context.Context, event Event) (string, e writeContext, cancel := detachedTimeout(ctx, emitter.timeout) defer cancel() result := emitter.database.WithContext(writeContext). + Session(&gorm.Session{SkipDefaultTransaction: true}). Model(&postgres.UsageOutbox{}). Where("event_id = ?", event.EventID). Updates(map[string]any{ @@ -108,26 +102,7 @@ func (emitter *DurableEmitter) Emit(ctx context.Context, event Event) (string, e if result.RowsAffected != 1 { return "", errors.New("usage outbox lifecycle record is missing") } - - entryID, publishErr := emitter.redis.Emit(ctx, event) - if publishErr != nil { - return "", publishErr - } - emitter.markPublished(event.EventID) - return entryID, nil -} - -func (emitter *DurableEmitter) markPublished(eventID string) { - ctx, cancel := context.WithTimeout(context.Background(), emitter.timeout) - defer cancel() - now := time.Now().UTC() - _ = emitter.database.WithContext(ctx). - Model(&postgres.UsageOutbox{}). - Where("event_id = ? AND state = ?", eventID, postgres.UsageOutboxReady). - Updates(map[string]any{ - "state": postgres.UsageOutboxPublished, - "published_at": now, - }).Error + return "", nil } // OutboxRelay republishes ready records and safely republishes published @@ -135,8 +110,8 @@ func (emitter *DurableEmitter) markPublished(eventID string) { type OutboxRelay struct { database *gorm.DB emitter *RedisEmitter + consumerGroup string batchSize int - timeout time.Duration pendingTimeout time.Duration republishAfter time.Duration } @@ -144,21 +119,24 @@ type OutboxRelay struct { func NewOutboxRelay( database *gorm.DB, emitter *RedisEmitter, - timeout time.Duration, + consumerGroup string, + batchSize int64, ) (*OutboxRelay, error) { switch { case database == nil: return nil, errors.New("usage outbox relay requires PostgreSQL") case emitter == nil: return nil, errors.New("usage outbox relay requires Redis") - case timeout <= 0: - return nil, errors.New("usage outbox relay timeout must be positive") + case consumerGroup == "": + return nil, errors.New("usage outbox relay requires a consumer group") + case batchSize <= 0 || batchSize > 1_000: + return nil, errors.New("usage outbox relay batch size is invalid") default: return &OutboxRelay{ database: database, emitter: emitter, - batchSize: defaultOutboxBatchSize, - timeout: timeout, + consumerGroup: consumerGroup, + batchSize: int(batchSize), pendingTimeout: 15 * time.Minute, republishAfter: defaultOutboxRepublishPeriod, }, nil @@ -177,46 +155,110 @@ func (relay *OutboxRelay) Publish(ctx context.Context) (int, error) { if _, err := relay.recoverPending(ctx); err != nil { return 0, err } + + records, err := relay.readyRecords(ctx) + if err != nil { + return 0, err + } + if len(records) > 0 { + return relay.publishRecords(ctx, records) + } + + caughtUp, err := relay.consumerCaughtUp(ctx) + if err != nil { + return 0, err + } + if !caughtUp { + return 0, nil + } + + records, err = relay.stalePublishedRecords(ctx) + if err != nil { + return 0, err + } + return relay.publishRecords(ctx, records) +} + +func (relay *OutboxRelay) readyRecords( + ctx context.Context, +) ([]postgres.UsageOutbox, error) { + var records []postgres.UsageOutbox + if err := relay.database.WithContext(ctx). + Where("state = ?", postgres.UsageOutboxReady). + Order("updated_at ASC"). + Limit(relay.batchSize). + Find(&records).Error; err != nil { + return nil, fmt.Errorf("read ready usage outbox: %w", err) + } + return records, nil +} + +func (relay *OutboxRelay) stalePublishedRecords( + ctx context.Context, +) ([]postgres.UsageOutbox, error) { republishBefore := time.Now().UTC().Add(-relay.republishAfter) var records []postgres.UsageOutbox if err := relay.database.WithContext(ctx). Where( - "state = ? OR (state = ? AND (published_at IS NULL OR published_at <= ?))", - postgres.UsageOutboxReady, + "state = ? AND (published_at IS NULL OR published_at <= ?)", postgres.UsageOutboxPublished, republishBefore, ). - Order("updated_at ASC"). + Order("published_at ASC"). Limit(relay.batchSize). Find(&records).Error; err != nil { - return 0, fmt.Errorf("read usage outbox: %w", err) + return nil, fmt.Errorf("read published usage outbox: %w", err) + } + return records, nil +} + +func (relay *OutboxRelay) consumerCaughtUp(ctx context.Context) (bool, error) { + groups, err := relay.emitter.client.XInfoGroups(ctx, relay.emitter.stream).Result() + if err != nil { + return false, fmt.Errorf("read usage consumer group: %w", err) + } + for _, group := range groups { + if group.Name == relay.consumerGroup { + return group.Pending == 0 && group.Lag == 0, nil + } } + return false, fmt.Errorf("usage consumer group %q was not found", relay.consumerGroup) +} - published := 0 +func (relay *OutboxRelay) publishRecords( + ctx context.Context, + records []postgres.UsageOutbox, +) (int, error) { + eventIDs := make([]string, 0, len(records)) for _, record := range records { if record.Payload == nil { - return published, fmt.Errorf("usage outbox event %s has no payload", record.EventID) + return 0, fmt.Errorf("usage outbox event %s has no payload", record.EventID) } event, err := Decode([]byte(*record.Payload)) if err != nil { - return published, fmt.Errorf("decode usage outbox event %s: %w", record.EventID, err) + return 0, fmt.Errorf("decode usage outbox event %s: %w", record.EventID, err) } if _, err := relay.emitter.Emit(ctx, event); err != nil { - return published, err + return 0, err } - now := time.Now().UTC() - if err := relay.database.WithContext(ctx). - Model(&postgres.UsageOutbox{}). - Where("event_id = ?", record.EventID). - Updates(map[string]any{ - "state": postgres.UsageOutboxPublished, - "published_at": now, - }).Error; err != nil { - return published, fmt.Errorf("mark usage outbox published: %w", err) - } - published++ + eventIDs = append(eventIDs, record.EventID) + } + if len(eventIDs) == 0 { + return 0, nil + } + + now := time.Now().UTC() + if err := relay.database.WithContext(ctx). + Session(&gorm.Session{SkipDefaultTransaction: true}). + Model(&postgres.UsageOutbox{}). + Where("event_id IN ?", eventIDs). + Updates(map[string]any{ + "state": postgres.UsageOutboxPublished, + "published_at": now, + }).Error; err != nil { + return 0, fmt.Errorf("mark usage outbox published: %w", err) } - return published, nil + return len(eventIDs), nil } func (relay *OutboxRelay) recoverPending(ctx context.Context) (int, error) { diff --git a/internal/usage/store.go b/internal/usage/store.go index 2eab8ce..7aabcca 100644 --- a/internal/usage/store.go +++ b/internal/usage/store.go @@ -20,6 +20,11 @@ type Store struct { now func() time.Time } +type storeEntry struct { + entryID string + event Event +} + func NewStore(database *gorm.DB, pricing *PricingCatalog) (*Store, error) { if database == nil { return nil, errors.New("usage store requires PostgreSQL") @@ -78,33 +83,61 @@ func (store *Store) ReloadManagedPricing(ctx context.Context) (bool, error) { } func (store *Store) Put(ctx context.Context, entryID string, event Event) (bool, error) { - if err := event.Validate(); err != nil { + stored, duplicates, err := store.putBatch( + ctx, + []storeEntry{{entryID: entryID, event: event}}, + ) + if err != nil { return false, err } - record := store.usageRecord(entryID, event, store.now().UTC()) - duplicate := false + return stored == 0 && duplicates == 1, nil +} + +func (store *Store) putBatch( + ctx context.Context, + entries []storeEntry, +) (int64, int64, error) { + if len(entries) == 0 { + return 0, 0, nil + } + processedAt := store.now().UTC() + records := make([]postgres.UsageEvent, 0, len(entries)) + eventIDs := make([]string, 0, len(entries)) + for _, entry := range entries { + if err := entry.event.Validate(); err != nil { + return 0, 0, err + } + records = append(records, store.usageRecord( + entry.entryID, + entry.event, + processedAt, + )) + eventIDs = append(eventIDs, entry.event.EventID) + } + + var stored int64 err := store.database.WithContext(ctx).Transaction(func(transaction *gorm.DB) error { result := transaction. Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "event_id"}}, DoNothing: true, }). - Create(&record) + Create(&records) if result.Error != nil { return result.Error } - duplicate = result.RowsAffected == 0 + stored = result.RowsAffected if err := transaction. - Where("event_id = ?", event.EventID). + Where("event_id IN ?", eventIDs). Delete(&postgres.UsageOutbox{}).Error; err != nil { return err } return nil }) if err != nil { - return false, err + return 0, 0, err } - return duplicate, nil + return stored, int64(len(entries)) - stored, nil } func (store *Store) usageRecord(entryID string, event Event, processedAt time.Time) postgres.UsageEvent { diff --git a/internal/usage/usage_test.go b/internal/usage/usage_test.go index 52987a9..22433f7 100644 --- a/internal/usage/usage_test.go +++ b/internal/usage/usage_test.go @@ -34,7 +34,7 @@ func TestCollectorFinalizesOnce(t *testing.T) { if err != nil { t.Fatalf("NewCollector() error = %v", err) } - collector.SetCacheStatus("bypass") + collector.SetCacheStatus("BYPASS") collector.SetRoute("primary", "upstream-model", 3, 1, 1) collector.ObserveResponse([]byte( `{"usage":{"prompt_tokens":11,"completion_tokens":4,"total_tokens":15}}`, @@ -51,6 +51,7 @@ func TestCollectorFinalizesOnce(t *testing.T) { t.Fatal("second Finalize() = true") } if event.LatencyMS != 1500 || + event.CacheStatus != "bypass" || event.Attempts != 3 || event.Retries != 1 || event.Fallbacks != 1 || @@ -513,7 +514,7 @@ func TestUsageRedisPostgresPipeline(t *testing.T) { t.Fatalf("first Worker.Run() error = %v", err) } - durable, err := NewDurableEmitter(database.ORM(), emitter, time.Second) + durable, err := NewDurableEmitter(database.ORM(), time.Second) if err != nil { t.Fatalf("NewDurableEmitter() error = %v", err) } @@ -530,10 +531,15 @@ func TestUsageRedisPostgresPipeline(t *testing.T) { if err != nil { t.Fatalf("DurableEmitter.Emit() error = %v", err) } - if err := client.XDel(ctx, settings.StreamKey, relayEntryID).Err(); err != nil { - t.Fatalf("delete published event before worker storage: %v", err) + if relayEntryID != "" { + t.Fatalf("DurableEmitter.Emit() entry ID = %q, want asynchronous relay", relayEntryID) } - relay, err := NewOutboxRelay(database.ORM(), emitter, time.Second) + relay, err := NewOutboxRelay( + database.ORM(), + emitter, + settings.Group, + settings.BatchSize, + ) if err != nil { t.Fatalf("NewOutboxRelay() error = %v", err) } @@ -542,7 +548,30 @@ func TestUsageRedisPostgresPipeline(t *testing.T) { t.Fatalf("OutboxRelay.Publish() published=%d error=%v", published, err) } if length, err := client.XLen(ctx, settings.StreamKey).Result(); err != nil || length != 1 { - t.Fatalf("republished stream length=%d error=%v", length, err) + t.Fatalf("relayed stream length=%d error=%v", length, err) + } + if published, err := relay.Publish(ctx); err != nil || published != 0 { + t.Fatalf("OutboxRelay.Publish(backlogged) published=%d error=%v", published, err) + } + messages, err := client.XReadGroup(ctx, &goredis.XReadGroupArgs{ + Group: settings.Group, + Consumer: settings.Consumer, + Streams: []string{settings.StreamKey, ">"}, + Count: 1, + }).Result() + if err != nil || len(messages) != 1 || len(messages[0].Messages) != 1 { + t.Fatalf("XReadGroup(relayed) streams=%d error=%v", len(messages), err) + } + lostEntryID := messages[0].Messages[0].ID + if _, err := client.TxPipelined(ctx, func(pipe goredis.Pipeliner) error { + pipe.XAck(ctx, settings.StreamKey, settings.Group, lostEntryID) + pipe.XDel(ctx, settings.StreamKey, lostEntryID) + return nil + }); err != nil { + t.Fatalf("remove relayed event before storage: %v", err) + } + if published, err := relay.Publish(ctx); err != nil || published != 1 { + t.Fatalf("OutboxRelay.Publish(lost) published=%d error=%v", published, err) } if duplicate, err := store.Put(ctx, "direct-relay", relayEvent); err != nil || duplicate { t.Fatalf("Put(relayed event) duplicate=%t error=%v", duplicate, err) diff --git a/internal/usage/worker.go b/internal/usage/worker.go index 463ff7b..835f918 100644 --- a/internal/usage/worker.go +++ b/internal/usage/worker.go @@ -260,49 +260,65 @@ func (worker *Worker) processBatch(ctx context.Context, messages []goredis.XMess ) defer cancel() + entries := make([]storeEntry, 0, len(messages)) + entryIDs := make([]string, 0, len(messages)) for _, message := range messages { - if err := worker.processMessage(batchContext, message); err != nil { - worker.stats.failed.Add(1) - slog.Error( - "usage worker entry failed", - "entry_id", message.ID, - "error", err, + payload, ok := streamString(message.Values["payload"]) + if !ok { + worker.recordEntryError( + message.ID, + worker.handlePoison(batchContext, message, "missing_payload"), ) + continue } + event, err := Decode([]byte(payload)) + if err != nil { + worker.recordEntryError( + message.ID, + worker.handlePoison(batchContext, message, "invalid_event"), + ) + continue + } + entries = append(entries, storeEntry{entryID: message.ID, event: event}) + entryIDs = append(entryIDs, message.ID) if batchContext.Err() != nil { return } } -} - -func (worker *Worker) processMessage(ctx context.Context, message goredis.XMessage) error { - payload, ok := streamString(message.Values["payload"]) - if !ok { - return worker.handlePoison(ctx, message, "missing_payload") + if len(entries) == 0 { + return } - event, err := Decode([]byte(payload)) + + stored, duplicates, err := worker.store.putBatch(batchContext, entries) if err != nil { - return worker.handlePoison(ctx, message, "invalid_event") + worker.stats.failed.Add(1) + slog.Error("usage worker batch failed", "entries", len(entries), "error", err) + return } + worker.stats.stored.Add(stored) + worker.stats.duplicates.Add(duplicates) - duplicate, err := worker.store.Put(ctx, message.ID, event) - if err != nil { - return err + if err := worker.acknowledge(batchContext, entryIDs); err != nil { + worker.stats.failed.Add(1) + slog.Error("usage worker acknowledgement failed", "entries", len(entryIDs), "error", err) } - if duplicate { - worker.stats.duplicates.Add(1) - } else { - worker.stats.stored.Add(1) +} + +func (worker *Worker) recordEntryError(entryID string, err error) { + if err == nil { + return } - _, err = worker.client.TxPipelined(ctx, func(pipe goredis.Pipeliner) error { - pipe.XAck(ctx, worker.config.StreamKey, worker.config.Group, message.ID) - pipe.XDel(ctx, worker.config.StreamKey, message.ID) + worker.stats.failed.Add(1) + slog.Error("usage worker entry failed", "entry_id", entryID, "error", err) +} + +func (worker *Worker) acknowledge(ctx context.Context, entryIDs []string) error { + _, err := worker.client.TxPipelined(ctx, func(pipe goredis.Pipeliner) error { + pipe.XAck(ctx, worker.config.StreamKey, worker.config.Group, entryIDs...) + pipe.XDel(ctx, worker.config.StreamKey, entryIDs...) return nil }) - if err != nil { - return err - } - return nil + return err } func (worker *Worker) handlePoison( diff --git a/test/threehost/collect-run-evidence.sh b/test/threehost/collect-run-evidence.sh index 86d2dd0..21db847 100644 --- a/test/threehost/collect-run-evidence.sh +++ b/test/threehost/collect-run-evidence.sh @@ -140,17 +140,17 @@ redis_command() { drain_deadline=$((SECONDS + drain_timeout_seconds)) drain_state=timeout -unfinished_outbox=-1 +remaining_outbox=-1 stream_length=-1 while ((SECONDS < drain_deadline)); do - unfinished_outbox=$(docker compose -f "$compose_file" exec -T postgres \ + remaining_outbox=$(docker compose -f "$compose_file" exec -T postgres \ psql \ -U "$postgres_user" \ -d "$postgres_database" \ -At \ - -c "SELECT count(*) FROM usage_outbox WHERE left(request_id, $prefix_length) = '$request_prefix' AND state <> 'published';") + -c "SELECT count(*) FROM usage_outbox WHERE left(request_id, $prefix_length) = '$request_prefix';") stream_length=$(redis_command --raw XLEN "$usage_stream") - if [[ "$unfinished_outbox" == "0" && "$stream_length" == "0" ]]; then + if [[ "$remaining_outbox" == "0" ]]; then drain_state=complete break fi @@ -159,7 +159,7 @@ done { printf 'state=%s\n' "$drain_state" printf 'timeout_seconds=%s\n' "$drain_timeout_seconds" - printf 'unfinished_outbox=%s\n' "$unfinished_outbox" + printf 'remaining_outbox=%s\n' "$remaining_outbox" printf 'stream_length=%s\n' "$stream_length" } >"$output_dir/usage-drain.txt" From 592564232b11d4aac7310ddd89d04fd13c94c73c Mon Sep 17 00:00:00 2001 From: tonghzhang <949674415@qq.com> Date: Tue, 28 Jul 2026 22:13:15 +0800 Subject: [PATCH 3/9] fix(usage): batch durable lifecycle writes --- cmd/model-velo/main.go | 7 + internal/usage/outbox.go | 241 +++++++++++++++++++++++++++++++---- internal/usage/usage_test.go | 74 +++++++++++ 3 files changed, 298 insertions(+), 24 deletions(-) diff --git a/cmd/model-velo/main.go b/cmd/model-velo/main.go index 714e632..ec0792a 100644 --- a/cmd/model-velo/main.go +++ b/cmd/model-velo/main.go @@ -117,6 +117,13 @@ func run() error { // 装配基础设施、启动 HTTP 服务并等待关闭。 if err != nil { return fmt.Errorf("configure durable usage emitter: %w", err) } + defer func() { + closeContext, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := usageEmitter.Close(closeContext); err != nil { + slog.Error("durable usage emitter shutdown failed", "error", err) + } + }() usagePricing, err := usage.NewPricingCatalog(startup.usage.Pricing) if err != nil { return fmt.Errorf("configure usage pricing: %w", err) diff --git a/internal/usage/outbox.go b/internal/usage/outbox.go index 3f3e30d..a20698c 100644 --- a/internal/usage/outbox.go +++ b/internal/usage/outbox.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "sync" "time" "gorm.io/gorm" @@ -14,6 +15,12 @@ import ( const defaultOutboxRepublishPeriod = 30 * time.Second +const ( + durableEmitterBatchSize = 100 + durableEmitterBatchWait = 5 * time.Millisecond + durableEmitterQueueSize = 4_096 +) + // PendingEvent is the durable minimum recorded before an authenticated request // reaches a provider. It intentionally excludes prompts and provider secrets. type PendingEvent struct { @@ -33,11 +40,24 @@ type LifecycleEmitter interface { Begin(context.Context, PendingEvent) error } -// DurableEmitter stores request lifecycle state in PostgreSQL. The worker owns -// Redis publication so online requests do not compete with replay traffic. +// DurableEmitter stores request lifecycle state in PostgreSQL. Concurrent +// requests share short, bounded batches so durability does not require two +// individual SQL round trips per request. type DurableEmitter struct { database *gorm.DB timeout time.Duration + writes chan durableWrite + stop chan struct{} + done chan struct{} + stateMu sync.RWMutex + closed bool + stopOnce sync.Once +} + +type durableWrite struct { + record postgres.UsageOutbox + ready bool + result chan error } func NewDurableEmitter( @@ -50,7 +70,15 @@ func NewDurableEmitter( case timeout <= 0: return nil, errors.New("durable usage emitter timeout must be positive") default: - return &DurableEmitter{database: database, timeout: timeout}, nil + emitter := &DurableEmitter{ + database: database, + timeout: timeout, + writes: make(chan durableWrite, durableEmitterQueueSize), + stop: make(chan struct{}), + done: make(chan struct{}), + } + go emitter.run() + return emitter, nil } } @@ -58,8 +86,6 @@ func (emitter *DurableEmitter) Begin(ctx context.Context, pending PendingEvent) if err := validatePendingEvent(pending); err != nil { return err } - writeContext, cancel := detachedTimeout(ctx, emitter.timeout) - defer cancel() record := postgres.UsageOutbox{ EventID: pending.EventID, RequestID: pending.RequestID, @@ -70,12 +96,8 @@ func (emitter *DurableEmitter) Begin(ctx context.Context, pending PendingEvent) State: postgres.UsageOutboxPending, StartedAt: pending.StartedAt.UTC(), } - result := emitter.database.WithContext(writeContext). - Session(&gorm.Session{SkipDefaultTransaction: true}). - Clauses(clause.OnConflict{Columns: []clause.Column{{Name: "event_id"}}, DoNothing: true}). - Create(&record) - if result.Error != nil { - return fmt.Errorf("record usage request lifecycle: %w", result.Error) + if err := emitter.enqueue(ctx, durableWrite{record: record}); err != nil { + return fmt.Errorf("record usage request lifecycle: %w", err) } return nil } @@ -85,24 +107,195 @@ func (emitter *DurableEmitter) Emit(ctx context.Context, event Event) (string, e if err != nil { return "", err } + payloadText := string(payload) + record := postgres.UsageOutbox{ + EventID: event.EventID, + RequestID: event.RequestID, + TenantID: event.TenantID, + APIKeyID: event.APIKeyID, + RequestedModel: event.RequestedModel, + Stream: event.Stream, + State: postgres.UsageOutboxReady, + Payload: &payloadText, + StartedAt: event.StartedAt.UTC(), + } + if err := emitter.enqueue(ctx, durableWrite{record: record, ready: true}); err != nil { + return "", fmt.Errorf("finalize usage outbox event: %w", err) + } + return "", nil +} + +func (emitter *DurableEmitter) Close(ctx context.Context) error { + emitter.stopOnce.Do(func() { + emitter.stateMu.Lock() + emitter.closed = true + close(emitter.stop) + emitter.stateMu.Unlock() + }) + if ctx == nil { + ctx = context.Background() + } + select { + case <-emitter.done: + return nil + case <-ctx.Done(): + return fmt.Errorf("close durable usage emitter: %w", ctx.Err()) + } +} + +func (emitter *DurableEmitter) enqueue( + ctx context.Context, + write durableWrite, +) error { writeContext, cancel := detachedTimeout(ctx, emitter.timeout) defer cancel() - result := emitter.database.WithContext(writeContext). + write.result = make(chan error, 1) + + emitter.stateMu.RLock() + if emitter.closed { + emitter.stateMu.RUnlock() + return errors.New("durable usage emitter is closed") + } + select { + case emitter.writes <- write: + emitter.stateMu.RUnlock() + case <-emitter.done: + emitter.stateMu.RUnlock() + return errors.New("durable usage emitter is closed") + case <-writeContext.Done(): + emitter.stateMu.RUnlock() + return writeContext.Err() + } + + select { + case err := <-write.result: + return err + case <-emitter.done: + return errors.New("durable usage emitter closed before write completed") + case <-writeContext.Done(): + return writeContext.Err() + } +} + +func (emitter *DurableEmitter) run() { + defer close(emitter.done) + for { + select { + case first := <-emitter.writes: + emitter.writeBatch(emitter.collect(first, false)) + case <-emitter.stop: + emitter.drain() + return + } + } +} + +func (emitter *DurableEmitter) collect( + first durableWrite, + stopping bool, +) []durableWrite { + batch := make([]durableWrite, 0, durableEmitterBatchSize) + batch = append(batch, first) + if stopping { + for len(batch) < durableEmitterBatchSize { + select { + case write := <-emitter.writes: + batch = append(batch, write) + default: + return batch + } + } + return batch + } + + timer := time.NewTimer(durableEmitterBatchWait) + defer timer.Stop() + for len(batch) < durableEmitterBatchSize { + select { + case write := <-emitter.writes: + batch = append(batch, write) + case <-timer.C: + return batch + case <-emitter.stop: + return batch + } + } + return batch +} + +func (emitter *DurableEmitter) drain() { + for { + select { + case first := <-emitter.writes: + emitter.writeBatch(emitter.collect(first, true)) + default: + return + } + } +} + +func (emitter *DurableEmitter) writeBatch(batch []durableWrite) { + pending := make([]postgres.UsageOutbox, 0, len(batch)) + ready := make([]postgres.UsageOutbox, 0, len(batch)) + for _, write := range batch { + if write.ready { + ready = append(ready, write.record) + continue + } + pending = append(pending, write.record) + } + + pendingErr := emitter.writePending(pending) + readyErr := emitter.writeReady(ready) + for _, write := range batch { + if write.ready { + write.result <- readyErr + continue + } + write.result <- pendingErr + } +} + +func (emitter *DurableEmitter) writePending( + records []postgres.UsageOutbox, +) error { + if len(records) == 0 { + return nil + } + ctx, cancel := context.WithTimeout(context.Background(), emitter.timeout) + defer cancel() + if err := emitter.database.WithContext(ctx). Session(&gorm.Session{SkipDefaultTransaction: true}). - Model(&postgres.UsageOutbox{}). - Where("event_id = ?", event.EventID). - Updates(map[string]any{ - "payload": string(payload), - "state": postgres.UsageOutboxReady, - "published_at": nil, - }) - if result.Error != nil { - return "", fmt.Errorf("finalize usage outbox event: %w", result.Error) + Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "event_id"}}, + DoNothing: true, + }). + CreateInBatches(&records, len(records)).Error; err != nil { + return err } - if result.RowsAffected != 1 { - return "", errors.New("usage outbox lifecycle record is missing") + return nil +} + +func (emitter *DurableEmitter) writeReady( + records []postgres.UsageOutbox, +) error { + if len(records) == 0 { + return nil } - return "", nil + ctx, cancel := context.WithTimeout(context.Background(), emitter.timeout) + defer cancel() + if err := emitter.database.WithContext(ctx). + Session(&gorm.Session{SkipDefaultTransaction: true}). + Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "event_id"}}, + DoUpdates: clause.AssignmentColumns([]string{ + "payload", "state", "published_at", "updated_at", + }), + }). + CreateInBatches(&records, len(records)).Error; err != nil { + return err + } + return nil } // OutboxRelay republishes ready records and safely republishes published diff --git a/internal/usage/usage_test.go b/internal/usage/usage_test.go index 22433f7..e626eeb 100644 --- a/internal/usage/usage_test.go +++ b/internal/usage/usage_test.go @@ -518,6 +518,80 @@ func TestUsageRedisPostgresPipeline(t *testing.T) { if err != nil { t.Fatalf("NewDurableEmitter() error = %v", err) } + defer func() { + closeContext, closeCancel := context.WithTimeout(context.Background(), 2*time.Second) + defer closeCancel() + if err := durable.Close(closeContext); err != nil { + t.Errorf("DurableEmitter.Close() error = %v", err) + } + }() + + const concurrentLifecycleCount = 200 + concurrentEvents := make([]Event, 0, concurrentLifecycleCount) + concurrentEventIDs := make([]string, 0, concurrentLifecycleCount) + for index := 0; index < concurrentLifecycleCount; index++ { + event := integrationEvent(t, "request-batched-"+strconv.Itoa(index)) + concurrentEvents = append(concurrentEvents, event) + concurrentEventIDs = append(concurrentEventIDs, event.EventID) + } + var lifecycleWait sync.WaitGroup + lifecycleErrors := make(chan error, concurrentLifecycleCount) + for _, event := range concurrentEvents { + lifecycleWait.Add(1) + go func(event Event) { + defer lifecycleWait.Done() + if err := durable.Begin(ctx, PendingEvent{ + EventID: event.EventID, RequestID: event.RequestID, + TenantID: event.TenantID, APIKeyID: event.APIKeyID, + RequestedModel: event.RequestedModel, Stream: event.Stream, + StartedAt: event.StartedAt, + }); err != nil { + lifecycleErrors <- err + return + } + if _, err := durable.Emit(ctx, event); err != nil { + lifecycleErrors <- err + } + }(event) + } + lifecycleWait.Wait() + close(lifecycleErrors) + for err := range lifecycleErrors { + t.Errorf("batched durable lifecycle error = %v", err) + } + var batchedReady int64 + if err := database.ORM().Model(&postgres.UsageOutbox{}). + Where("event_id IN ? AND state = ?", concurrentEventIDs, postgres.UsageOutboxReady). + Count(&batchedReady).Error; err != nil || batchedReady != concurrentLifecycleCount { + t.Fatalf( + "batched ready lifecycles = %d, want %d, error = %v", + batchedReady, concurrentLifecycleCount, err, + ) + } + if err := database.ORM(). + Where("event_id IN ?", concurrentEventIDs). + Delete(&postgres.UsageOutbox{}).Error; err != nil { + t.Fatalf("delete batched usage lifecycles: %v", err) + } + + finalOnlyEvent := integrationEvent(t, "request-final-without-begin") + if _, err := durable.Emit(ctx, finalOnlyEvent); err != nil { + t.Fatalf("DurableEmitter.Emit(without Begin) error = %v", err) + } + var finalOnly postgres.UsageOutbox + if err := database.ORM(). + Where("event_id = ?", finalOnlyEvent.EventID). + First(&finalOnly).Error; err != nil || + finalOnly.State != postgres.UsageOutboxReady || + finalOnly.Payload == nil { + t.Fatalf("final-only outbox row = %#v, error = %v", finalOnly, err) + } + if err := database.ORM(). + Where("event_id = ?", finalOnlyEvent.EventID). + Delete(&postgres.UsageOutbox{}).Error; err != nil { + t.Fatalf("delete final-only usage lifecycle: %v", err) + } + relayEvent := integrationEvent(t, "request-outbox-republish") if err := durable.Begin(ctx, PendingEvent{ EventID: relayEvent.EventID, RequestID: relayEvent.RequestID, From a16991c881be235007a0ef68f0d7b2188fed9740 Mon Sep 17 00:00:00 2001 From: tonghzhang <949674415@qq.com> Date: Tue, 28 Jul 2026 22:13:34 +0800 Subject: [PATCH 4/9] test(perf): add final reliability evidence gates --- .gitignore | 2 + test/fakeupstream/main.go | 2 +- test/fakeupstream/server.go | 88 ++++++-- test/fakeupstream/server_test.go | 36 ++++ test/k6/common.js | 4 + test/threehost/benchmark.env.example | 6 + test/threehost/collect-run-evidence.sh | 2 +- test/threehost/report-template.html | 101 ++++++--- test/threehost/run-complete-client.sh | 79 ++++++- test/threehost/run-usage-chaos.sh | 282 +++++++++++++++++++++++++ test/threehost/summarize.py | 211 +++++++++++++++++- 11 files changed, 746 insertions(+), 67 deletions(-) create mode 100644 test/threehost/run-usage-chaos.sh diff --git a/.gitignore b/.gitignore index 9353a82..e31fcd0 100644 --- a/.gitignore +++ b/.gitignore @@ -11,6 +11,8 @@ *.out coverage.* *.bench +__pycache__/ +*.pyc /test-results/ /test/threehost/client.env /test/threehost/benchmark.env diff --git a/test/fakeupstream/main.go b/test/fakeupstream/main.go index 7f6e029..bc5b416 100644 --- a/test/fakeupstream/main.go +++ b/test/fakeupstream/main.go @@ -95,7 +95,7 @@ func run(config commandConfig) error { "name", upstream.providerName, "scenario", - upstream.scenarioOverride, + upstream.forcedScenario(), ) select { diff --git a/test/fakeupstream/server.go b/test/fakeupstream/server.go index 21330c3..400e0b9 100644 --- a/test/fakeupstream/server.go +++ b/test/fakeupstream/server.go @@ -144,8 +144,9 @@ var scenarioCatalog = map[string]scenario{ } type upstreamServer struct { - providerName string - scenarioOverride string + providerName string + scenarioMu sync.RWMutex + scenario string attemptsMu sync.Mutex attempts map[string]attemptState @@ -169,9 +170,11 @@ type attemptState struct { } type scenarioStats struct { - Requests int64 `json:"requests"` - Errors int64 `json:"errors"` - Streams int64 `json:"streams"` + Requests int64 + Errors int64 + Streams int64 + FirstRequestAt time.Time + LastRequestAt time.Time } func newUpstreamServer(providerName, scenarioOverride string) (*upstreamServer, error) { @@ -189,10 +192,10 @@ func newUpstreamServer(providerName, scenarioOverride string) (*upstreamServer, } } server := &upstreamServer{ - providerName: providerName, - scenarioOverride: scenarioOverride, - attempts: map[string]attemptState{}, - scenarioStats: map[string]scenarioStats{}, + providerName: providerName, + scenario: scenarioOverride, + attempts: map[string]attemptState{}, + scenarioStats: map[string]scenarioStats{}, } server.statsStartedAt.Store(time.Now().UnixNano()) return server, nil @@ -204,6 +207,7 @@ func (s *upstreamServer) handler() http.Handler { mux.HandleFunc("GET /__admin/scenarios", s.handleScenarios) mux.HandleFunc("GET /__admin/stats", s.handleStats) mux.HandleFunc("POST /__admin/reset", s.handleReset) + mux.HandleFunc("POST /__admin/scenario", s.handleScenario) mux.HandleFunc("POST /v1/chat/completions", s.handleChatCompletions) mux.HandleFunc("POST /chat/completions", s.handleChatCompletions) return mux @@ -213,7 +217,7 @@ func (s *upstreamServer) handleHealth(w http.ResponseWriter, _ *http.Request) { writeJSON(w, http.StatusOK, healthResponse{ Status: "ok", Provider: s.providerName, - ScenarioOverride: s.scenarioOverride, + ScenarioOverride: s.forcedScenario(), }) } @@ -234,7 +238,7 @@ func (s *upstreamServer) handleScenarios(w http.ResponseWriter, _ *http.Request) } writeJSON(w, http.StatusOK, scenariosResponse{ Default: defaultScenarioName, - Override: s.scenarioOverride, + Override: s.forcedScenario(), Scenarios: available, }) } @@ -250,10 +254,12 @@ func (s *upstreamServer) handleStats(w http.ResponseWriter, _ *http.Request) { for _, name := range names { stats := s.scenarioStats[name] scenarios = append(scenarios, scenarioStatsEntry{ - Name: name, - Requests: stats.Requests, - Errors: stats.Errors, - Streams: stats.Streams, + Name: name, + Requests: stats.Requests, + Errors: stats.Errors, + Streams: stats.Streams, + FirstRequestAt: stats.FirstRequestAt, + LastRequestAt: stats.LastRequestAt, }) } s.statsMu.Unlock() @@ -289,6 +295,31 @@ func (s *upstreamServer) handleReset(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusNoContent) } +func (s *upstreamServer) handleScenario(w http.ResponseWriter, r *http.Request) { + decoder := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<10)) + decoder.DisallowUnknownFields() + var request scenarioOverrideRequest + if err := decoder.Decode(&request); err != nil { + writeUpstreamError(w, http.StatusBadRequest, "invalid scenario override") + return + } + if err := decoder.Decode(&struct{}{}); !errors.Is(err, io.EOF) { + writeUpstreamError(w, http.StatusBadRequest, "invalid scenario override") + return + } + request.Scenario = strings.TrimSpace(request.Scenario) + if request.Scenario != "" { + if _, exists := scenarioCatalog[request.Scenario]; !exists { + writeUpstreamError(w, http.StatusBadRequest, "unknown scenario override") + return + } + } + s.scenarioMu.Lock() + s.scenario = request.Scenario + s.scenarioMu.Unlock() + w.WriteHeader(http.StatusNoContent) +} + func (s *upstreamServer) handleChatCompletions(w http.ResponseWriter, r *http.Request) { s.beginRequest() defer s.finishRequest() @@ -388,6 +419,11 @@ func (s *upstreamServer) recordScenario(name string, stream bool) { s.statsMu.Lock() stats := s.scenarioStats[name] stats.Requests++ + now := time.Now().UTC() + if stats.FirstRequestAt.IsZero() { + stats.FirstRequestAt = now + } + stats.LastRequestAt = now if stream { stats.Streams++ s.streams.Add(1) @@ -409,7 +445,7 @@ func (s *upstreamServer) recordError(name string) { } func (s *upstreamServer) selectScenario(model string) (scenario, error) { - name := s.scenarioOverride + name := s.forcedScenario() if name == "" && strings.HasPrefix(model, "mock/") { name = model } @@ -423,6 +459,12 @@ func (s *upstreamServer) selectScenario(model string) (scenario, error) { return selected, nil } +func (s *upstreamServer) forcedScenario() string { + s.scenarioMu.RLock() + defer s.scenarioMu.RUnlock() + return s.scenario +} + func (s *upstreamServer) beginAttempt( requestID string, selected scenario, @@ -836,6 +878,10 @@ type healthResponse struct { ScenarioOverride string `json:"scenario_override,omitempty"` } +type scenarioOverrideRequest struct { + Scenario string `json:"scenario"` +} + type scenarioDescription struct { Name string `json:"name"` Description string `json:"description"` @@ -848,10 +894,12 @@ type scenariosResponse struct { } type scenarioStatsEntry struct { - Name string `json:"name"` - Requests int64 `json:"requests"` - Errors int64 `json:"errors"` - Streams int64 `json:"streams"` + Name string `json:"name"` + Requests int64 `json:"requests"` + Errors int64 `json:"errors"` + Streams int64 `json:"streams"` + FirstRequestAt time.Time `json:"first_request_at"` + LastRequestAt time.Time `json:"last_request_at"` } type upstreamStatsResponse struct { diff --git a/test/fakeupstream/server_test.go b/test/fakeupstream/server_test.go index 4d6db73..b08bb6a 100644 --- a/test/fakeupstream/server_test.go +++ b/test/fakeupstream/server_test.go @@ -208,6 +208,42 @@ func TestUpstreamServerLoadScenariosAndStats(t *testing.T) { } } +func TestUpstreamServerScenarioOverrideCanRecover(t *testing.T) { + upstream, err := newUpstreamServer("recovering-provider", "mock/error-503") + if err != nil { + t.Fatalf("new upstream server: %v", err) + } + testServer := httptest.NewServer(upstream.handler()) + defer testServer.Close() + + failed := postChat(t, testServer.URL, "before-recovery", "mock/instant", false) + _, _ = io.Copy(io.Discard, failed.Body) + _ = failed.Body.Close() + if failed.StatusCode != http.StatusServiceUnavailable { + t.Fatalf("status before recovery = %d, want 503", failed.StatusCode) + } + + response, err := http.Post( + testServer.URL+"/__admin/scenario", + "application/json", + strings.NewReader(`{"scenario":"mock/instant"}`), + ) + if err != nil { + t.Fatalf("set scenario override: %v", err) + } + _, _ = io.Copy(io.Discard, response.Body) + _ = response.Body.Close() + if response.StatusCode != http.StatusNoContent { + t.Fatalf("scenario update status = %d, want 204", response.StatusCode) + } + + recovered := postChat(t, testServer.URL, "after-recovery", "mock/error-503", false) + defer recovered.Body.Close() + if recovered.StatusCode != http.StatusOK { + t.Fatalf("status after recovery = %d, want 200", recovered.StatusCode) + } +} + func requestIDWithFailureOutcome(fail bool) string { for index := range 10_000 { requestID := "failure-outcome-" + strconv.Itoa(index) diff --git a/test/k6/common.js b/test/k6/common.js index bc12cdb..e42506c 100644 --- a/test/k6/common.js +++ b/test/k6/common.js @@ -8,6 +8,7 @@ export const responses200 = new Counter('chat_responses_200'); export const responses429 = new Counter('chat_responses_429'); export const responses5xx = new Counter('chat_responses_5xx'); export const responsesOther = new Counter('chat_responses_other'); +export const gatewayChatRequests = new Counter('gateway_chat_requests'); export const requestTimeout = __ENV.REQUEST_TIMEOUT || '20s'; export const expected200 = http.expectedStatuses(200); @@ -125,6 +126,9 @@ export function sendChat({ stream, }); + if (target === 'gateway') { + gatewayChatRequests.add(1, { model, stream: String(stream) }); + } return http.post(chatURL(targetURL), body, { headers, responseCallback, diff --git a/test/threehost/benchmark.env.example b/test/threehost/benchmark.env.example index aec3696..cba2516 100644 --- a/test/threehost/benchmark.env.example +++ b/test/threehost/benchmark.env.example @@ -1,6 +1,8 @@ # Private addresses reachable from the client host. GATEWAY_URL=http://10.0.0.20:8080 UPSTREAM_URL=http://10.0.0.30:9000 +UPSTREAM_FAIL_URL=http://10.0.0.30:9001 +UPSTREAM_FALLBACK_URL=http://10.0.0.30:9002 MODEL_VELO_API_KEY=replace-with-model-velo-api-key # Use the same run ID on the client, gateway, and upstream hosts. @@ -45,6 +47,9 @@ PROFILE_MAX_VUS=2048 # Fault and queue-overload cases. FAULT_RATE=100 FAULT_DURATION=2m +FAULT_RECOVERY_RATE=100 +FAULT_RECOVERY_FAILURE_DURATION=45s +FAULT_RECOVERY_HEALTHY_DURATION=45s QUEUE_RATE=1000 QUEUE_DURATION=2m QUEUE_PRE_ALLOCATED_VUS=2048 @@ -71,6 +76,7 @@ RUN_CACHE=true RUN_RAMP=true RUN_BURST=true RUN_FAULT=true +RUN_FAULT_RECOVERY=true RUN_QUEUE_OVERLOAD=true RUN_ENDURANCE=true RUN_RELIABILITY=true diff --git a/test/threehost/collect-run-evidence.sh b/test/threehost/collect-run-evidence.sh index 21db847..a0b5115 100644 --- a/test/threehost/collect-run-evidence.sh +++ b/test/threehost/collect-run-evidence.sh @@ -189,7 +189,7 @@ docker compose -f "$compose_file" exec -T postgres \ -U "$postgres_user" \ -d "$postgres_database" \ --csv \ - -c "SELECT requested_model, status, stream, cache_status, error_category, error_code, count(*) AS events, sum(attempts) AS attempts, sum(retries) AS retries, sum(fallbacks) AS fallbacks, round(avg(latency_ms)::numeric, 2) AS latency_avg_ms, percentile_cont(0.95) WITHIN GROUP (ORDER BY latency_ms) AS latency_p95_ms, round(avg(first_token_ms)::numeric, 2) AS first_token_avg_ms FROM usage_events WHERE left(request_id, $prefix_length) = '$request_prefix' GROUP BY requested_model, status, stream, cache_status, error_category, error_code ORDER BY requested_model, status, stream, cache_status, error_category, error_code;" \ + -c "SELECT requested_model, status, stream, cache_status, error_category, error_code, count(*) AS events, sum(attempts) AS attempts, sum(retries) AS retries, sum(fallbacks) AS fallbacks, round(avg(latency_ms)::numeric, 2) AS latency_avg_ms, percentile_cont(0.95) WITHIN GROUP (ORDER BY latency_ms) AS latency_p95_ms, percentile_cont(0.99) WITHIN GROUP (ORDER BY latency_ms) AS latency_p99_ms, round(avg(first_token_ms)::numeric, 2) AS first_token_avg_ms, percentile_cont(0.95) WITHIN GROUP (ORDER BY first_token_ms) AS first_token_p95_ms, percentile_cont(0.99) WITHIN GROUP (ORDER BY first_token_ms) AS first_token_p99_ms FROM usage_events WHERE left(request_id, $prefix_length) = '$request_prefix' GROUP BY requested_model, status, stream, cache_status, error_category, error_code ORDER BY requested_model, status, stream, cache_status, error_category, error_code;" \ >"$output_dir/usage-diagnostics.csv" { diff --git a/test/threehost/report-template.html b/test/threehost/report-template.html index ddb9c41..69c74df 100644 --- a/test/threehost/report-template.html +++ b/test/threehost/report-template.html @@ -459,12 +459,11 @@ 03 / SSE 流式 04 / 可靠性与耐久 05 / 热路径诊断 - 06 / 客户端资源 + 06 / 资源与镜像 07 / 证据完整性 - 08 / 其他网关对比 - 09 / 环境与口径 - 10 / 全部 Case - 11 / 原始产物 + 08 / 环境与口径 + 09 / 全部 Case + 10 / 原始产物
    本文件自包含,可离线打开。报告含私网地址与运行元数据,不建议直接公开部署。
    @@ -517,12 +516,12 @@

    当前最可辩护的结论

    闭环网关吞吐

    -

    c=1…256,三次运行取中位数。吞吐非单调,c=128 后约 1.25k RPS;这不是稳定 SLO。

    +

    每档并发按重复运行取中位数;闭环峰值用于定位吞吐拐点,不等同于生产 SLO。

    闭环 P99

    -

    并发上升后排队长尾快速放大:c=256 的 P99 比直连多约 364 ms。

    +

    同时展示 P50 与 P99,便于识别并发上升后是否出现排队和尾延迟放大。

    @@ -550,7 +549,7 @@

    直连基线 vs 网关增量

    03 / server-sent events

    SSE 回放与逐 Chunk 观测

    -

    500 个流式请求、并发 20、三次重复。独立 stream loader 记录 headers、首事件、首内容、总耗时和 chunk 间隔,不只看 k6 的普通 HTTP duration。

    +

    独立 stream loader 记录 headers、首事件、首内容、总耗时和 chunk 间隔;请求数、并发和重复次数均取自本次运行元数据。

    @@ -560,16 +559,14 @@

    首内容与总完成时间

    怎么解读

    -

    首内容 P50 增量只有约 5.37 ms,说明正常路径的转发、SSE 解析与首 chunk 提交成本较小;但首内容 P99 增量约 50 ms,总耗时 P99 也多约 52 ms,尾部抖动值得继续用网关 CPU profile 和 per-stage timing 分解。

    -

    inter‑chunk P99 只增加约 0.60 ms,说明流一旦开始,持续转发没有明显积累性停顿。三轮 3,000 个直连/网关流请求全部完整成功。

    -
    边界已经测到。可靠性脚本同时覆盖首 chunk 前非法事件返回 502,以及首 chunk 已提交后断流不再 Fallback;这两条是 SSE 网关最容易写错的状态边界。
    +
    -
    04 / failure semantics

    缓存、故障、队列与 30 分钟耐久

    +
    04 / failure semantics

    缓存、故障恢复、队列与耐久

    这一组不是单纯追求成功率:error10 和 queue overload 会有意制造可接受错误。报告将“测试断言成功”和“业务 HTTP 2xx”分开。

    @@ -579,10 +576,8 @@

    17 项可靠性断言

      -

      三条最重要的诊断

      -

      重试放大:error10 的上游调用约为客户端请求的 1.20×。由于错误由 request ID 确定,重试不会“治愈”这 10%,只会增加上游压力。

      -

      Queue 边界:1,000 RPS 过载 case 中约一半请求成功,上游最大 active 恰为 256;被拒请求没有继续打到上游,边界有效,但 P99 已接近 3 秒。

      -

      耐久错误:500 RPS、30 分钟共 900,001 请求,518 个 5xx,业务成功率 99.942%。它通过 99.9%,但没有达到 99.99%。

      +

      本次可靠性诊断

      +
      @@ -619,8 +614,8 @@

      该 case 的网关错误码

      -
      06 / load generator

      客户端资源:峰值属于谁

      -

      客户端有 5,835 个约 1 秒采样。服务器文件也已收到,但它们分别止于主运行前或始于主运行后,必须先做时间窗口审计。

      +
      06 / runtime envelope

      客户端、网关容器与镜像

      +

      资源数据只使用与 case 时间窗重叠的采样;同时展示客户端负载生成器、网关容器内存和最终镜像体积。

      @@ -660,7 +655,7 @@

      证据审计结论

      -
      +