diff --git a/.github/workflows/pr-path-guard.yml b/.github/workflows/pr-path-guard.yml index 720140332..30d5635db 100644 --- a/.github/workflows/pr-path-guard.yml +++ b/.github/workflows/pr-path-guard.yml @@ -23,6 +23,7 @@ jobs: env: TWP_NC_MAINTAINERS: ${{ vars.TWP_NC_MAINTAINERS }} ACTOR: ${{ github.actor }} + PR_AUTHOR: ${{ github.event.pull_request.user.login }} shell: bash run: | set -euo pipefail @@ -31,7 +32,9 @@ jobs: is_maintainer=false for m in "${MAINTAINERS[@]}"; do m="$(echo "$m" | xargs)" - if [[ "$ACTOR" == "$m" ]]; then + # A cloud agent push sets actor to cursor[bot] on a maintainer-owned PR. + # The PR author is the account that opened it. + if [[ "$ACTOR" == "$m" || "$PR_AUTHOR" == "$m" ]]; then is_maintainer=true break fi diff --git a/.github/workflows/rps-ab.yml b/.github/workflows/rps-ab.yml new file mode 100644 index 000000000..a0b7644b2 --- /dev/null +++ b/.github/workflows/rps-ab.yml @@ -0,0 +1,193 @@ +# Same-job paired A/B for Titanium RPS work (workflow_dispatch only; Linux only). +# +# Why: two rps-saturation dispatches land on two different VMs and runner variance (~5%) is larger than +# the gains left on the Linux gap list. This job checks out a baseline and a candidate product ref on ONE +# ubuntu-latest VM, builds the SAME probe harness (tools/RpsLoadProbe from the dispatch ref) against each, +# and runs them alternately baseline, candidate, baseline, candidate ... for >= 5 pairs with a cool-down +# between runs. Paired deltas, a 95% confidence interval and a keep gate are written to the job summary. +# +# Dispatch the workflow on the harness ref (usually the candidate branch): +# gh workflow run rps-ab.yml --ref \ +# -f baseline_ref=develop -f candidate_ref= \ +# -f mode=compare-ceiling -f arm_contains=twp-reverse-http1 -f pairs=5 +# +# Only product code (everything outside tools/RpsLoadProbe) differs between the two sides. +# Peers (nginx / HAProxy / Envoy) are not installed here; keep arm_contains on Titanium (twp-) or bare- arms. + +name: RPS paired A/B + +on: + workflow_dispatch: + inputs: + baseline_ref: + description: 'Baseline product ref (branch, tag or SHA)' + required: true + default: develop + candidate_ref: + description: 'Candidate product ref (branch, tag or SHA)' + required: false + default: '' + mode: + description: 'Probe ramp mode (see rps-saturation.yml for the list)' + required: true + default: compare-ceiling + arm_contains: + description: 'Arm-name substring filter (applied to the arms of the mode)' + required: true + default: 'twp-reverse-http1' + concurrency: + description: 'Concurrency (one value; paired runs use a fixed step)' + required: true + default: '32' + duration_sec: + description: 'Measure seconds' + required: true + default: '8' + pairs: + description: 'Number of baseline/candidate pairs (>= 5 for a keep decision)' + required: true + default: '5' + cooldown_sec: + description: 'Idle seconds between consecutive runs' + required: true + default: '15' + order: + description: 'alternate = A,B,A,B ; counterbalance = A,B,B,A,A,B (cancels linear drift)' + required: true + default: alternate + type: choice + options: + - alternate + - counterbalance + probe_args: + description: 'Extra probe arguments, space separated; later flags win (e.g. "--response-bytes 65536 --warmup-sec 3 --arm-shard 1/3")' + required: false + default: '' + +permissions: + contents: read + +jobs: + ab: + runs-on: ubuntu-latest + timeout-minutes: 300 + env: + AB_BASELINE_REF: ${{ inputs.baseline_ref }} + AB_CANDIDATE_REF: ${{ inputs.candidate_ref }} + AB_MODE: ${{ inputs.mode }} + AB_ARM_CONTAINS: ${{ inputs.arm_contains }} + AB_CONCURRENCY: ${{ inputs.concurrency }} + AB_DURATION: ${{ inputs.duration_sec }} + AB_PAIRS: ${{ inputs.pairs }} + AB_COOLDOWN: ${{ inputs.cooldown_sec }} + AB_ORDER: ${{ inputs.order }} + AB_PROBE_ARGS: ${{ inputs.probe_args }} + steps: + - uses: actions/checkout@v7 + with: + fetch-depth: 0 + + - uses: actions/setup-dotnet@v6 + with: + dotnet-version: | + 10.0.x + + - name: Log runner shape + run: | + set -euo pipefail + echo "os=$(uname -s) $(uname -r)" + echo "nproc=$(nproc)" + lscpu | sed -n '1,20p' + + - name: Install libmsquic (HTTP/3 arms) + run: | + set -euo pipefail + . /etc/os-release + curl --fail --silent --show-error --location --proto '=https' --tlsv1.2 \ + "https://packages.microsoft.com/config/${ID}/${VERSION_ID}/packages-microsoft-prod.deb" \ + -o packages-microsoft-prod.deb + sudo dpkg -i packages-microsoft-prod.deb + rm -f packages-microsoft-prod.deb + sudo find /etc/apt -type f \( -iname '*chrome*' -o -iname '*google*chrome*' \) -print -delete || true + sudo apt-get update || true + sudo apt-get install -y libmsquic + + - name: Check out baseline and candidate worktrees + run: | + set -euo pipefail + git fetch --no-tags origin '+refs/heads/*:refs/remotes/origin/*' || true + resolve() { + git rev-parse --verify --quiet "$1^{commit}" \ + || git rev-parse --verify --quiet "origin/$1^{commit}" + } + base_sha="$(resolve "$AB_BASELINE_REF")" || { echo "cannot resolve baseline_ref=$AB_BASELINE_REF" >&2; exit 1; } + cand_ref="${AB_CANDIDATE_REF:-$GITHUB_SHA}" + cand_sha="$(resolve "$cand_ref")" || { echo "cannot resolve candidate_ref=$cand_ref" >&2; exit 1; } + echo "baseline $AB_BASELINE_REF -> $base_sha" + echo "candidate $cand_ref -> $cand_sha" + echo "harness $GITHUB_SHA (tools/RpsLoadProbe overlaid on both sides)" + git worktree add --detach "$RUNNER_TEMP/ab/baseline" "$base_sha" + git worktree add --detach "$RUNNER_TEMP/ab/candidate" "$cand_sha" + { + echo "AB_BASELINE_SHA=$base_sha" + echo "AB_CANDIDATE_SHA=$cand_sha" + } >> "$GITHUB_ENV" + + - name: Build the identical probe against each product ref + run: | + set -euo pipefail + for side in baseline candidate; do + wt="$RUNNER_TEMP/ab/$side" + rm -rf "$wt/tools/RpsLoadProbe" + cp -r tools/RpsLoadProbe "$wt/tools/RpsLoadProbe" + dotnet build -c Release "$wt/tools/RpsLoadProbe/RpsLoadProbe.csproj" --warnaserror + mkdir -p "$RUNNER_TEMP/ab-bin" + cp -r "$wt/tools/RpsLoadProbe/bin/Release/net10.0" "$RUNNER_TEMP/ab-bin/$side" + done + ls -la "$RUNNER_TEMP/ab-bin/baseline/RpsLoadProbe" "$RUNNER_TEMP/ab-bin/candidate/RpsLoadProbe" + + - name: Paired A/B ramp + run: | + set -euo pipefail + ulimit -n 65535 || true + extra=() + if [ -n "${AB_PROBE_ARGS:-}" ]; then + # shellcheck disable=SC2206 + for a in $AB_PROBE_ARGS; do extra+=(--probe-arg "$a"); done + fi + python3 tools/RpsLoadProbe/rps-ab-run.py \ + --baseline-probe "$RUNNER_TEMP/ab-bin/baseline/RpsLoadProbe" \ + --candidate-probe "$RUNNER_TEMP/ab-bin/candidate/RpsLoadProbe" \ + --mode "$AB_MODE" --arm-contains "$AB_ARM_CONTAINS" \ + --concurrency "$AB_CONCURRENCY" --duration-sec "$AB_DURATION" \ + --pairs "$AB_PAIRS" --cooldown-sec "$AB_COOLDOWN" --order "$AB_ORDER" \ + "${extra[@]}" --out "$RUNNER_TEMP/ab-out" + + - name: Analyze paired deltas + if: always() + run: | + set -euo pipefail + test -f "$RUNNER_TEMP/ab-out/manifest.json" || exit 0 + { + echo "baseline \`$AB_BASELINE_REF\` = \`$AB_BASELINE_SHA\` " + echo "candidate \`${AB_CANDIDATE_REF:-$GITHUB_SHA}\` = \`$AB_CANDIDATE_SHA\` " + echo "harness (identical on both sides) = \`$GITHUB_SHA\`" + echo + } > "$RUNNER_TEMP/ab-out/summary.md" + python3 tools/RpsLoadProbe/rps-ab-analyze.py "$RUNNER_TEMP/ab-out" --markdown "$RUNNER_TEMP/ab-out/analysis.md" > /dev/null + cat "$RUNNER_TEMP/ab-out/analysis.md" >> "$RUNNER_TEMP/ab-out/summary.md" + cat "$RUNNER_TEMP/ab-out/summary.md" >> "$GITHUB_STEP_SUMMARY" + cat "$RUNNER_TEMP/ab-out/summary.md" + + - name: Upload A/B artifacts + if: always() + uses: actions/upload-artifact@v7 + with: + name: rps-ab-${{ github.run_id }} + path: | + ${{ runner.temp }}/ab-out/**/*.csv + ${{ runner.temp }}/ab-out/**/*.tsv + ${{ runner.temp }}/ab-out/**/run.log + ${{ runner.temp }}/ab-out/manifest.json + ${{ runner.temp }}/ab-out/summary.md + if-no-files-found: error diff --git a/.github/workflows/rps-profile.yml b/.github/workflows/rps-profile.yml new file mode 100644 index 000000000..6089307bc --- /dev/null +++ b/.github/workflows/rps-profile.yml @@ -0,0 +1,165 @@ +# Linux CPU / allocation / perf capture for a few RPS arms (workflow_dispatch only; diagnosis, not publishable RPS). +# +# One arm at a time, fixed concurrency, ~60 s measure per arm. Per arm the capture is +# - perf record -g on the whole proxy process tree (kernel + OpenSSL + runtime + JIT frames named +# via DOTNET_PerfMapEnabled and libcoreclr / libclrjit symbols from the Microsoft symbol server) +# - dotnet-trace dotnet-sampled-thread-time (managed stacks) +# - dotnet-trace gc-verbose -> GCAllocationTick -> bytes allocated per request by type +# - dotnet-counters System.Runtime (alloc rate, GC counts, lock contentions, thread-pool queue) +# Native peers (nginx / HAProxy) get the perf window only: that is the native floor for the same wire. +# summarize-profile.py turns the captures into a markdown table (job summary + artifact). +# +# gh workflow run rps-profile.yml --ref \ +# -f arms='reverse-http1-tls,bare-reverse-http1-tls,nginx-reverse-http1-tls,haproxy-reverse-http1-tls' +# +# Peers use the distro nginx / HAProxy builds (H1 + TLS terminate only); use rps-saturation.yml for +# publishable peer numbers. + +name: RPS profile (Linux) + +on: + workflow_dispatch: + inputs: + arms: + description: 'Comma-separated single-arm probe modes (e.g. reverse-http1,bare-reverse-http1,nginx-reverse-http1)' + required: true + default: 'reverse-http1-tls,bare-reverse-http1-tls,nginx-reverse-http1-tls,haproxy-reverse-http1-tls' + concurrency: + description: 'Fixed concurrency for every captured arm' + required: true + default: '32' + capture_sec: + description: 'Seconds per capture window (perf, sampled-thread-time, gc-verbose); measure step = 3 windows + slack' + required: true + default: '20' + probe_args: + description: 'Extra probe arguments, space separated (e.g. "--response-bytes 65536")' + required: false + default: '' + +permissions: + contents: read + +jobs: + profile: + runs-on: ubuntu-latest + timeout-minutes: 120 + env: + PF_ARMS: ${{ inputs.arms }} + PF_CONCURRENCY: ${{ inputs.concurrency }} + PF_CAPTURE_SEC: ${{ inputs.capture_sec }} + PF_PROBE_ARGS: ${{ inputs.probe_args }} + steps: + - uses: actions/checkout@v7 + + - uses: actions/setup-dotnet@v6 + with: + dotnet-version: | + 10.0.x + + - name: Log runner shape + run: | + set -euo pipefail + echo "os=$(uname -s) $(uname -r)" + echo "nproc=$(nproc)" + lscpu | sed -n '1,20p' + + - name: Install profiling tools, nginx, HAProxy and libmsquic + run: | + set -euo pipefail + sudo find /etc/apt -type f \( -iname '*chrome*' -o -iname '*google*chrome*' \) -print -delete || true + sudo apt-get update || true + sudo apt-get install -y linux-tools-common "linux-tools-$(uname -r)" || sudo apt-get install -y linux-tools-generic || true + # linux-tools-generic ships perf under /usr/lib/linux-tools// ; expose one that runs. + perf_bin="$(command -v perf || true)" + if [ -z "$perf_bin" ] || ! "$perf_bin" --version >/dev/null 2>&1; then + cand="$(ls -d /usr/lib/linux-tools/*/perf 2>/dev/null | tail -n 1 || true)" + [ -n "$cand" ] && sudo ln -sf "$cand" /usr/local/bin/perf + fi + perf --version + sudo apt-get install -y nginx haproxy || true + sudo systemctl stop nginx haproxy || true + . /etc/os-release + curl --fail --silent --show-error --location --proto '=https' --tlsv1.2 \ + "https://packages.microsoft.com/config/${ID}/${VERSION_ID}/packages-microsoft-prod.deb" \ + -o packages-microsoft-prod.deb + sudo dpkg -i packages-microsoft-prod.deb + rm -f packages-microsoft-prod.deb + sudo apt-get update || true + sudo apt-get install -y libmsquic || true + dotnet tool install -g dotnet-trace + dotnet tool install -g dotnet-counters + dotnet tool install -g dotnet-symbol + echo "$HOME/.dotnet/tools" >> "$GITHUB_PATH" + nginx -v || true + haproxy -v | head -n 1 || true + + - name: Install runtime native symbols (libcoreclr, libclrjit) for perf + run: | + set -euo pipefail + export PATH="$PATH:$HOME/.dotnet/tools" + rt="$(dirname "$(dotnet --list-runtimes | awk '/Microsoft.NETCore.App 10\./ {gsub(/\[|\]/,"",$3); print $3 "/" $2 "/x"}' | tail -n 1)")" + echo "runtime dir: $rt" + mkdir -p "$RUNNER_TEMP/sym" + dotnet-symbol --symbols --output "$RUNNER_TEMP/sym" "$rt/libcoreclr.so" "$rt/libclrjit.so" + for f in libcoreclr libclrjit; do + bid="$(readelf -n "$rt/$f.so" | awk '/Build ID/ {print $3}')" + d="/usr/lib/debug/.build-id/${bid:0:2}" + sudo mkdir -p "$d" + sudo cp "$RUNNER_TEMP/sym/$f.so.dbg" "$d/${bid:2}.debug" + echo "$f build-id $bid" + done + + - name: Build probe and allocation-tick reader + run: | + set -euo pipefail + dotnet build -c Release tools/RpsLoadProbe/RpsLoadProbe.csproj --warnaserror + dotnet build -c Release tools/RpsAllocTicks/RpsAllocTicks.csproj + + - name: Capture arms + run: | + set -euo pipefail + export PATH="$PATH:$HOME/.dotnet/tools" + ulimit -n 65535 || true + export RPS_ALLOC_TICKS="$PWD/tools/RpsAllocTicks/bin/Release/net10.0/RpsAllocTicks" + extra=() + if [ -n "${PF_PROBE_ARGS:-}" ]; then + # shellcheck disable=SC2206 + for a in $PF_PROBE_ARGS; do extra+=(--probe-arg "$a"); done + fi + IFS=',' read -ra arms <<< "$PF_ARMS" + for arm in "${arms[@]}"; do + arm="$(echo "$arm" | xargs)" + [ -z "$arm" ] && continue + echo "::group::profile $arm" + bash tools/RpsLoadProbe/profile-arm.sh --probe "$PWD/tools/RpsLoadProbe/bin/Release/net10.0/RpsLoadProbe" \ + --mode "$arm" --out "$RUNNER_TEMP/profile/$arm" --concurrency "$PF_CONCURRENCY" \ + --capture-sec "$PF_CAPTURE_SEC" "${extra[@]}" || echo "::warning::profile of $arm failed" + echo "::endgroup::" + sleep 10 + done + + - name: Summarise + if: always() + run: | + set -euo pipefail + test -d "$RUNNER_TEMP/profile" || exit 0 + python3 tools/RpsLoadProbe/summarize-profile.py "$RUNNER_TEMP/profile" --cpus "$(nproc)" \ + --markdown "$RUNNER_TEMP/profile/summary.md" > /dev/null + cat "$RUNNER_TEMP/profile/summary.md" >> "$GITHUB_STEP_SUMMARY" + + - name: Upload captures + if: always() + uses: actions/upload-artifact@v7 + with: + name: rps-profile-${{ github.run_id }} + path: | + ${{ runner.temp }}/profile/summary.md + ${{ runner.temp }}/profile/**/*.txt + ${{ runner.temp }}/profile/**/*.csv + ${{ runner.temp }}/profile/**/*.log + ${{ runner.temp }}/profile/**/*.tsv + ${{ runner.temp }}/profile/**/*.nettrace + ${{ runner.temp }}/profile/**/*.speedscope.json + ${{ runner.temp }}/profile/**/perf.data.gz + if-no-files-found: error diff --git a/.github/workflows/rps-saturation.yml b/.github/workflows/rps-saturation.yml index 9ae1d92cf..2c897e0f9 100644 --- a/.github/workflows/rps-saturation.yml +++ b/.github/workflows/rps-saturation.yml @@ -18,6 +18,48 @@ on: branches: [beta, stable] # Push editions are gated by .NET / rps-publish-gate (Ubuntu only). Do not duplicate # 3-OS compare-editions on every beta/stable merge. + # Called by rps-suite.yml: one job per wiki row. Inputs mirror workflow_dispatch but are + # plain strings (workflow_call does not support choice inputs). + workflow_call: + inputs: + mode: + description: 'Ramp mode' + required: true + type: string + concurrency: + description: 'Comma-separated concurrency steps' + required: false + type: string + default: '8,16,32,64' + warmup_sec: + description: 'Warmup seconds per step' + required: false + type: string + default: '2' + duration_sec: + description: 'Measure seconds per step' + required: false + type: string + default: '8' + repeats: + description: 'Full arm sequence repeats (median peaks)' + required: false + type: string + default: '3' + runner_os: + description: 'Single runner OS for this leg' + required: true + type: string + arm_shard: + description: 'all, i/n, or one comparison-group key' + required: false + type: string + default: all + arm_contains: + description: 'Optional arm-name substring filter' + required: false + type: string + default: '' workflow_dispatch: inputs: mode: @@ -90,7 +132,6 @@ on: - envoy-reverse-http3-to-http2 - envoy-reverse-http3-to-http3 - compare-editions - - compare-cross-version - compare-ceiling - compare-bodies - compare-post @@ -195,7 +236,7 @@ on: - windows-latest - macos-15 arm_shard: - description: 'Comparison-group partition (all = no split; i/n = exclusive 1-based shard of wiki rows / Client×Origin wires, not individual arms)' + description: 'Comparison-group partition (all = no split; i/n = exclusive 1-based shard of wiki rows; or one group key such as h1c-h1c or h3-h3-body64k = exactly one wiki row, see --print-groups)' required: true default: all arm_contains: @@ -212,13 +253,13 @@ jobs: permissions: contents: read env: - RPS_MODE: ${{ github.event_name == 'workflow_dispatch' && inputs.mode || 'compare-spot' }} - RPS_CONCURRENCY: ${{ github.event_name == 'workflow_dispatch' && inputs.concurrency || '8,16,32,64' }} - RPS_WARMUP_SEC: ${{ github.event_name == 'workflow_dispatch' && inputs.warmup_sec || 2 }} - RPS_DURATION_SEC: ${{ github.event_name == 'workflow_dispatch' && inputs.duration_sec || 8 }} - RPS_REPEATS: ${{ github.event_name == 'workflow_dispatch' && inputs.repeats || 1 }} - RPS_ARM_SHARD: ${{ github.event_name == 'workflow_dispatch' && inputs.arm_shard || 'all' }} - RPS_ARM_CONTAINS: ${{ github.event_name == 'workflow_dispatch' && inputs.arm_contains || '' }} + RPS_MODE: ${{ (github.event_name == 'workflow_dispatch' || github.event_name == 'workflow_call') && inputs.mode || 'compare-spot' }} + RPS_CONCURRENCY: ${{ (github.event_name == 'workflow_dispatch' || github.event_name == 'workflow_call') && inputs.concurrency || '8,16,32,64' }} + RPS_WARMUP_SEC: ${{ (github.event_name == 'workflow_dispatch' || github.event_name == 'workflow_call') && inputs.warmup_sec || 2 }} + RPS_DURATION_SEC: ${{ (github.event_name == 'workflow_dispatch' || github.event_name == 'workflow_call') && inputs.duration_sec || 8 }} + RPS_REPEATS: ${{ (github.event_name == 'workflow_dispatch' || github.event_name == 'workflow_call') && inputs.repeats || 1 }} + RPS_ARM_SHARD: ${{ (github.event_name == 'workflow_dispatch' || github.event_name == 'workflow_call') && inputs.arm_shard || 'all' }} + RPS_ARM_CONTAINS: ${{ (github.event_name == 'workflow_dispatch' || github.event_name == 'workflow_call') && inputs.arm_contains || '' }} strategy: fail-fast: false matrix: @@ -229,7 +270,7 @@ jobs: # GitHub-hosted runners hard-cap at 360 minutes — do not set timeout above that. # Wiki-grade compare-product / bodies / arch use arm_shard i/n so each job stays under 345m ramp. # compare-product-smoke / haproxy-smoke / envoy-smoke: 120m; compare-editions: 90m. - timeout-minutes: ${{ github.event_name == 'pull_request' && 45 || (github.event_name == 'push' && 90 || (inputs.mode == 'compare-editions' && 90 || ((inputs.mode == 'compare-product-smoke' || inputs.mode == 'compare-haproxy-smoke' || inputs.mode == 'compare-envoy-smoke') && 120 || 360))) }} + timeout-minutes: ${{ github.event_name == 'pull_request' && 45 || (github.event_name == 'push' && 90 || (inputs.mode == 'compare-editions' && 90 || ((inputs.mode == 'compare-product-smoke' || inputs.mode == 'compare-haproxy-smoke' || inputs.mode == 'compare-envoy-smoke') && 120 || ((inputs.arm_shard != 'all' && !contains(inputs.arm_shard, '/')) && 75 || 360)))) }} steps: - uses: actions/checkout@v7 @@ -835,7 +876,11 @@ jobs: # Pass workflow_dispatch inputs via env (not ${{ }} in the script) to avoid # githubactions:S7630 script-injection findings on user-controlled values. - name: Run saturation ramp - timeout-minutes: ${{ github.event_name == 'pull_request' && 40 || (github.event_name == 'push' && 60 || (inputs.mode == 'compare-editions' && 60 || ((inputs.mode == 'compare-product-smoke' || inputs.mode == 'compare-haproxy-smoke' || inputs.mode == 'compare-envoy-smoke') && 110 || 345))) }} + id: ramp + # Same cap on every OS (macOS included): the 180m macOS-only cap was hit by the full + # saturation ramp (run 37110567901). Leave headroom under the 360m job timeout so a true + # hang still uploads a partial CSV. + timeout-minutes: ${{ github.event_name == 'pull_request' && 40 || (github.event_name == 'push' && 60 || (inputs.mode == 'compare-editions' && 60 || ((inputs.mode == 'compare-product-smoke' || inputs.mode == 'compare-haproxy-smoke' || inputs.mode == 'compare-envoy-smoke') && 110 || ((inputs.arm_shard != 'all' && !contains(inputs.arm_shard, '/')) && 60 || 345)))) }} shell: pwsh env: RPS_MODE: ${{ env.RPS_MODE }} @@ -869,6 +914,22 @@ jobs: @containsArgs ` @hapArgs + - name: Mark incomplete ramp + if: always() && steps.ramp.outcome != 'success' + shell: pwsh + env: + RAMP_OUTCOME: ${{ steps.ramp.outcome }} + run: | + $dir = 'tools/RpsLoadProbe/results' + New-Item -ItemType Directory -Force -Path $dir | Out-Null + @( + "ramp_outcome=$env:RAMP_OUTCOME" + "marked_at_utc=$((Get-Date).ToUniversalTime().ToString('o'))" + "os=$env:RUNNER_OS" + "arm_shard=$env:RPS_ARM_SHARD" + ) | Set-Content -LiteralPath (Join-Path $dir 'INCOMPLETE.txt') -Encoding utf8 + Write-Host "Wrote INCOMPLETE.txt (ramp outcome=$env:RAMP_OUTCOME)" + - name: Sanitize artifact shard label if: always() shell: bash @@ -877,10 +938,24 @@ jobs: raw="${RPS_ARM_SHARD:-all}" echo "RPS_ARTIFACT_SHARD=$(printf '%s' "$raw" | tr '/' '-')" >> "$GITHUB_ENV" + # Every arm this host resolved must have rows in every repeat. Gate scripts skip + # pairs whose arms are absent, so a wedged arm would otherwise leave the leg green. + - name: Assert ramp completeness + if: steps.ramp.outcome == 'success' + id: completeness + shell: pwsh + run: | + pwsh tools/RpsLoadProbe/assert-ramp-complete.ps1 ` + -ResultsDir tools/RpsLoadProbe/results ` + -Mode $env:RPS_MODE ` + -ArmShard $env:RPS_ARM_SHARD ` + -Repeats ([int]$env:RPS_REPEATS) + # Post-ramp gates: hard-fail on every event (including workflow_dispatch). # fail-fast: false keeps sibling OS jobs running; Upload CSV uses if: always(). # Sharded CSVs skip pairs whose arms are not in this artifact (union paste validates completeness). - name: Validate edition gates + id: gate-editions if: env.RPS_MODE == 'compare-editions' shell: pwsh run: | @@ -890,16 +965,18 @@ jobs: pwsh tools/RpsLoadProbe/validate-edition-gates.ps1 -CsvPath $csv.FullName - name: Validate compare-product gates + id: gate-product if: env.RPS_MODE == 'compare-product' || env.RPS_MODE == 'compare-product-smoke' shell: pwsh run: | $csv = Get-ChildItem tools/RpsLoadProbe/results -Filter 'rps-ramp-*.csv' | Sort-Object LastWriteTime -Descending | Select-Object -First 1 if (-not $csv) { throw 'No CSV found for compare-product gate validation' } - # Same floors on all OS: MITM Lite÷Reverse >= 0.40, Full÷Reverse >= 0.40, reverse TWP÷YARP >= 0.60. No nginx gate. + # Same floors on all OS: MITM Lite÷Reverse >= 0.25, Full÷Reverse >= 0.25, reverse TWP÷closest peer >= 0.50. pwsh tools/RpsLoadProbe/validate-compare-product-gates.ps1 -CsvPath $csv.FullName - name: Validate heavier YARP gates (bodies) + id: gate-bodies if: env.RPS_MODE == 'compare-bodies' shell: pwsh run: | @@ -910,6 +987,7 @@ jobs: pwsh tools/RpsLoadProbe/validate-heavier-yarp-gates.ps1 -CsvPath $csv.FullName -NameSuffix body256k - name: Validate heavier YARP gates (post) + id: gate-post if: env.RPS_MODE == 'compare-post' shell: pwsh run: | @@ -919,6 +997,7 @@ jobs: pwsh tools/RpsLoadProbe/validate-heavier-yarp-gates.ps1 -CsvPath $csv.FullName -NameSuffix post64k - name: Validate lossy hard-bar gates + id: gate-lossy if: env.RPS_MODE == 'compare-lossy' shell: pwsh run: | @@ -928,6 +1007,7 @@ jobs: pwsh tools/RpsLoadProbe/validate-lossy-arch-gates.ps1 -CsvPath $csv.FullName -Suite lossy - name: Validate arch hard-bar gates + id: gate-arch if: env.RPS_MODE == 'compare-arch' shell: pwsh run: | @@ -937,6 +1017,7 @@ jobs: pwsh tools/RpsLoadProbe/validate-lossy-arch-gates.ps1 -CsvPath $csv.FullName -Suite arch - name: Validate gRPC YARP gates + id: gate-grpc if: env.RPS_MODE == 'compare-grpc' shell: pwsh run: | @@ -945,22 +1026,45 @@ jobs: if (-not $csv) { throw 'No CSV found for gRPC YARP gate validation' } pwsh tools/RpsLoadProbe/validate-grpc-yarp-gates.ps1 -CsvPath $csv.FullName - - name: Validate cross-version gates - if: env.RPS_MODE == 'compare-cross-version' + # One status file per leg, overwritten on re-run, so the suite aggregator reads the + # last attempt rather than a stale incomplete marker from an earlier attempt. + - name: Write leg status + if: always() shell: pwsh + env: + RAMP_OUTCOME: ${{ steps.ramp.outcome }} + COMPLETENESS_OUTCOME: ${{ steps.completeness.outcome }} + GATE_EDITIONS: ${{ steps.gate-editions.outcome }} + GATE_PRODUCT: ${{ steps.gate-product.outcome }} + GATE_BODIES: ${{ steps.gate-bodies.outcome }} + GATE_POST: ${{ steps.gate-post.outcome }} + GATE_LOSSY: ${{ steps.gate-lossy.outcome }} + GATE_ARCH: ${{ steps.gate-arch.outcome }} + GATE_GRPC: ${{ steps.gate-grpc.outcome }} run: | - $csv = Get-ChildItem tools/RpsLoadProbe/results -Filter 'rps-ramp-*.csv' | - Sort-Object LastWriteTime -Descending | Select-Object -First 1 - if (-not $csv) { throw 'No CSV found for cross-version gate validation' } - $baseline = if ($env:RUNNER_OS -eq 'Windows') { - 'tools/RpsLoadProbe/results/baseline-6.0-win.csv' - } else { - 'tools/RpsLoadProbe/results/baseline-6.0-linux.csv' - } - if (-not (Test-Path $baseline)) { throw "Missing baseline: $baseline" } - pwsh tools/RpsLoadProbe/validate-cross-version.ps1 ` - -BaselineCsv $baseline ` - -CurrentCsv $csv.FullName + $dir = 'tools/RpsLoadProbe/results' + New-Item -ItemType Directory -Force -Path $dir | Out-Null + $gateOutcomes = @($env:GATE_EDITIONS, $env:GATE_PRODUCT, $env:GATE_BODIES, $env:GATE_POST, $env:GATE_LOSSY, $env:GATE_ARCH, $env:GATE_GRPC) + $gateOutcome = if ($gateOutcomes -contains 'failure') { 'failure' } + elseif ($env:RAMP_OUTCOME -eq 'success' -and $env:COMPLETENESS_OUTCOME -eq 'success') { 'success' } + else { 'skipped' } + [ordered]@{ + ramp_outcome = $env:RAMP_OUTCOME + completeness_outcome = $env:COMPLETENESS_OUTCOME + gate_outcome = $gateOutcome + attempt = [int]$env:GITHUB_RUN_ATTEMPT + arm_shard = $env:RPS_ARM_SHARD + os = $env:RUNNER_OS + } | ConvertTo-Json | Set-Content -LiteralPath (Join-Path $dir 'leg-status.json') -Encoding utf8 + + - name: Upload leg status + if: always() + uses: actions/upload-artifact@v7 + with: + name: rps-status-${{ matrix.os }}-shard-${{ env.RPS_ARTIFACT_SHARD }} + path: tools/RpsLoadProbe/results/leg-status.json + overwrite: true + if-no-files-found: warn - name: Upload CSV artifact if: always() @@ -968,5 +1072,18 @@ jobs: with: # Include arm_shard so parallel shard dispatches do not collide (slash → dash via env sanitize). name: rps-csv-${{ matrix.os }}-shard-${{ env.RPS_ARTIFACT_SHARD }} - path: tools/RpsLoadProbe/results/rps-ramp-*.csv + # A re-run attempt of the same run must replace the first attempt's artifact. + overwrite: true + path: | + tools/RpsLoadProbe/results/rps-ramp-*.csv + tools/RpsLoadProbe/results/rps-ramp-*.tls.tsv if-no-files-found: error + + - name: Upload incomplete marker + if: always() && hashFiles('tools/RpsLoadProbe/results/INCOMPLETE.txt') != '' + uses: actions/upload-artifact@v7 + with: + name: rps-incomplete-${{ matrix.os }}-shard-${{ env.RPS_ARTIFACT_SHARD }} + overwrite: true + path: tools/RpsLoadProbe/results/INCOMPLETE.txt + if-no-files-found: ignore diff --git a/.github/workflows/rps-suite.yml b/.github/workflows/rps-suite.yml new file mode 100644 index 000000000..6ed53cd2e --- /dev/null +++ b/.github/workflows/rps-suite.yml @@ -0,0 +1,194 @@ +# One job per wiki row for a single ramp mode, on all three OS. +# Dispatch one run per mode (compare-product, compare-bodies, ...) so each table has one +# Source URL and the paste scripts, which look arms up by name, never mix modes. +# compare-saturation and compare-editions stay one job per OS (key "all"): saturation Block A +# compares proxy arms against origin-direct from the same job. +# +# Free-plan account limits: 20 concurrent jobs, 5 of them macOS. This workflow caps each OS +# matrix itself; runs of different modes still queue together against the account limit. + +name: RPS suite + +on: + workflow_dispatch: + inputs: + mode: + description: 'Ramp mode (one suite run per mode)' + required: true + default: compare-product + type: choice + options: + - compare-product + - compare-bodies + - compare-post + - compare-lossy + - compare-tls-cost + - compare-arch + - compare-grpc + - compare-ws-h1tls + - compare-ws-h2 + - compare-saturation + - compare-editions + repeats: + description: 'Repeats per arm (median peaks)' + required: true + default: '3' + max_parallel_mac: + description: 'Concurrent macOS legs (Free plan allows 5)' + required: true + default: '5' + max_parallel_other: + description: 'Concurrent Windows or Linux legs' + required: true + default: '15' + +permissions: + contents: read + actions: read + +jobs: + prep: + name: prep ${{ inputs.mode }} + runs-on: ubuntu-latest + outputs: + keys: ${{ steps.groups.outputs.keys }} + steps: + - uses: actions/checkout@v7 + + - name: Setup .NET + uses: actions/setup-dotnet@v6 + with: + dotnet-version: | + 10.0.x + + - name: Resolve comparison groups + id: groups + shell: pwsh + env: + SUITE_MODE: ${{ inputs.mode }} + run: | + $ErrorActionPreference = 'Stop' + # These modes compare arms that must share a job, so they stay unsharded. + $unsharded = @('compare-saturation', 'compare-editions') + if ($unsharded -contains $env:SUITE_MODE) { + $keys = @('all') + } + else { + dotnet build -c Release tools/RpsLoadProbe/RpsLoadProbe.csproj --nologo -v q | Out-Host + if ($LASTEXITCODE -ne 0) { throw 'Probe build failed' } + $lines = & dotnet run --project tools/RpsLoadProbe/RpsLoadProbe.csproj --no-build -c Release -- ` + --ramp --mode $env:SUITE_MODE --print-groups 2>$null + $groups = @($lines | Where-Object { $_ -match '^\S+\t\d+$' } | ForEach-Object { + $p = $_ -split "`t" + [pscustomobject]@{ Key = $p[0]; Arms = [int]$p[1] } + }) + if ($groups.Count -eq 0) { throw "No comparison groups for $($env:SUITE_MODE)" } + if ($groups.Count -gt 256) { throw "$($groups.Count) groups exceeds the 256-job matrix limit" } + # Longest rows first so the slowest work starts in the first wave. + $keys = @($groups | Sort-Object Arms -Descending | Select-Object -ExpandProperty Key) + } + Write-Host ("{0}: {1} group(s): {2}" -f $env:SUITE_MODE, $keys.Count, ($keys -join ', ')) + # Piping a one-element array into ConvertTo-Json unwraps it to a JSON string. + # fromJSON then is not an array, and the OS matrices are skipped (saturation, + # editions, and any mode with a single comparison group). + $keysJson = '[' + (($keys | ForEach-Object { ConvertTo-Json -InputObject ([string]$_) }) -join ',') + ']' + $roundTrip = @($keysJson | ConvertFrom-Json) + if ($roundTrip.Count -ne $keys.Count) { throw "keys output must be a JSON array of $($keys.Count): $keysJson" } + "keys=$keysJson" >> $env:GITHUB_OUTPUT + + linux: + name: ${{ inputs.mode }} / ${{ matrix.key }} / linux + needs: prep + strategy: + fail-fast: false + max-parallel: ${{ fromJSON(inputs.max_parallel_other) }} + matrix: + key: ${{ fromJSON(needs.prep.outputs.keys) }} + uses: ./.github/workflows/rps-saturation.yml + with: + mode: ${{ inputs.mode }} + runner_os: ubuntu-latest + arm_shard: ${{ matrix.key }} + repeats: ${{ inputs.repeats }} + + windows: + name: ${{ inputs.mode }} / ${{ matrix.key }} / windows + needs: prep + strategy: + fail-fast: false + max-parallel: ${{ fromJSON(inputs.max_parallel_other) }} + matrix: + key: ${{ fromJSON(needs.prep.outputs.keys) }} + uses: ./.github/workflows/rps-saturation.yml + with: + mode: ${{ inputs.mode }} + runner_os: windows-latest + arm_shard: ${{ matrix.key }} + repeats: ${{ inputs.repeats }} + + macos: + name: ${{ inputs.mode }} / ${{ matrix.key }} / macos + needs: prep + strategy: + fail-fast: false + max-parallel: ${{ fromJSON(inputs.max_parallel_mac) }} + matrix: + key: ${{ fromJSON(needs.prep.outputs.keys) }} + uses: ./.github/workflows/rps-saturation.yml + with: + mode: ${{ inputs.mode }} + runner_os: macos-15 + arm_shard: ${{ matrix.key }} + repeats: ${{ inputs.repeats }} + + aggregate: + name: aggregate ${{ inputs.mode }} + needs: [prep, linux, windows, macos] + if: always() + runs-on: ubuntu-latest + steps: + - name: Check every row on every OS + shell: pwsh + env: + SUITE_MODE: ${{ inputs.mode }} + KEYS_JSON: ${{ needs.prep.outputs.keys }} + LINUX_RESULT: ${{ needs.linux.result }} + WINDOWS_RESULT: ${{ needs.windows.result }} + MACOS_RESULT: ${{ needs.macos.result }} + run: | + $ErrorActionPreference = 'Stop' + $keys = @($env:KEYS_JSON | ConvertFrom-Json) + $arts = gh api "repos/$env:GITHUB_REPOSITORY/actions/runs/$env:GITHUB_RUN_ID/artifacts?per_page=100" --paginate --jq '.artifacts[] | select(.name | startswith("rps-status-")) | .name' + $statusNames = @($arts | Where-Object { $_ }) + $osLabels = [ordered]@{ linux = 'ubuntu-latest'; windows = 'windows-latest'; macos = 'macos-15' } + $missing = @() + foreach ($key in $keys) { + $label = ($key -replace '/', '-') + foreach ($os in $osLabels.Values) { + $name = "rps-status-$os-shard-$label" + if ($statusNames -notcontains $name) { $missing += $name } + } + } + $summary = @( + "### RPS suite: $env:SUITE_MODE", + "", + "| OS | Result |", + "|---|---|", + "| linux | $env:LINUX_RESULT |", + "| windows | $env:WINDOWS_RESULT |", + "| macos | $env:MACOS_RESULT |", + "", + "Expected rows: $($keys.Count) x 3 OS.", + "Status artifacts found: $($statusNames.Count)." + ) + if ($missing) { + $summary += "" + $summary += "Missing status artifacts:" + $summary += ($missing | ForEach-Object { "- $_" }) + } + $summary -join "`n" | Out-File -FilePath $env:GITHUB_STEP_SUMMARY -Encoding utf8 + $osFailed = @($env:LINUX_RESULT, $env:WINDOWS_RESULT, $env:MACOS_RESULT) | Where-Object { $_ -ne 'success' } + if ($missing -or $osFailed) { + throw "Suite incomplete: $($missing.Count) status artifact(s) missing, OS results: linux=$env:LINUX_RESULT windows=$env:WINDOWS_RESULT macos=$env:MACOS_RESULT" + } + Write-Host "Suite complete: $($keys.Count) rows x 3 OS" diff --git a/src/Titanium.Cli/Titanium.Cli.csproj b/src/Titanium.Cli/Titanium.Cli.csproj index 36675dd65..77df06a18 100644 --- a/src/Titanium.Cli/Titanium.Cli.csproj +++ b/src/Titanium.Cli/Titanium.Cli.csproj @@ -7,7 +7,7 @@ latest enable false - 7.0.15 + 7.0.16 Jehonathan Thomas Titanium Web Proxy CLI (titanium / twp). MIT diff --git a/src/Titanium.Inspector/Services/AvaloniaStatusNotifier.cs b/src/Titanium.Inspector/Services/AvaloniaStatusNotifier.cs index f97ab5938..35e27168e 100644 --- a/src/Titanium.Inspector/Services/AvaloniaStatusNotifier.cs +++ b/src/Titanium.Inspector/Services/AvaloniaStatusNotifier.cs @@ -1,3 +1,4 @@ +using Avalonia.Controls; using Avalonia.Controls.Notifications; namespace Titanium.Inspector.Services; @@ -18,7 +19,7 @@ public void Show(string message, StatusSeverity severity) } var manager = _manager(); - if (manager is null) + if (manager is null || !CanShowWindowToast(manager)) { return; } @@ -48,4 +49,24 @@ public void Show(string message, StatusSeverity severity) manager.Show(new Notification(title, message, type, duration)); } + + /// + /// is async void: it posts the card, then + /// closes it after Task.Delay. Headless drains that post in + /// ResetForUnitTests before the animation clock exists, and the delay + /// continuation runs after the dispatch clears the sync context, so Close + /// hits the thread pool and crashes the test host. The status bar still updates. + /// + private static bool CanShowWindowToast(WindowNotificationManager manager) + { + var top = TopLevel.GetTopLevel(manager); + var platform = (top as Window)?.PlatformImpl ?? top?.PlatformImpl; + var name = platform?.GetType().FullName; + if (string.IsNullOrEmpty(name)) + { + return false; + } + + return !name.Contains("Avalonia.Headless", StringComparison.Ordinal); + } } diff --git a/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs b/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs index bc4a0b901..0b0ddeaab 100644 --- a/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs +++ b/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs @@ -15,9 +15,14 @@ public sealed class SessionBodyDiskCache : IDisposable private const string LegacySessionFileSearchPattern = "*.json"; private readonly string _rootDirectory; - private readonly string _runDirectory; + private string _runDirectory; + private int _runGeneration; private long _maxBytes; private readonly object _gate = new(); + private readonly object _writeGate = new(); + private readonly object _cleanupGate = new(); + private readonly Queue _cleanupPaths = new(); + private Task _cleanupTask = Task.CompletedTask; /// Absolute file path → tracked size / time / session id. private readonly Dictionary _index = new( StringComparer.OrdinalIgnoreCase); @@ -41,16 +46,37 @@ internal SessionBodyDiskCache( RebuildIndexAndEnforceBudget(); } - /// Updates disk budget; returns current-run session ids whose files were deleted. + /// Updates disk budget. Index updates return immediately; file deletes run in the background. public IReadOnlyList UpdateLimits(long maxBytes, TimeSpan maxAge) { _ = maxAge; - lock (_gate) + return ScheduleBudgetPrune(maxBytes); + } + + /// Generation of the current run folder. Writes captured under an older generation are dropped. + internal int CurrentGeneration => Volatile.Read(ref _runGeneration); + + /// Thread that last drained queued file deletes. Zero until the first cleanup runs. + internal int LastCleanupThreadId { get; private set; } + + /// Waits until queued directory and file deletes have finished. + public async Task FlushCleanupAsync() + { + while (true) { - _maxBytes = maxBytes > 0 ? maxBytes : _maxBytes; - } + Task task; + lock (_cleanupGate) + { + if (_cleanupPaths.Count == 0 && _cleanupTask.IsCompleted) + { + return; + } - return EnforceDiskBudget(); + task = _cleanupTask; + } + + await task.ConfigureAwait(false); + } } /// @@ -78,9 +104,29 @@ public string PathFor(long sessionId) => public bool FileExists(long sessionId) => File.Exists(PathFor(sessionId)); /// Writes a single-entry HAR into the current run folder; may prune older files under budget. - public IReadOnlyList Write(SessionSnapshot snapshot) + public IReadOnlyList Write(SessionSnapshot snapshot) => + Write(snapshot, expectedGeneration: null); + + /// + /// Writes when still matches. A clear that rotated the run + /// folder rejects the write so an in-flight spill cannot recreate the abandoned session id. + /// + internal IReadOnlyList Write(SessionSnapshot snapshot, int? expectedGeneration) + { + lock (_writeGate) + { + ObjectDisposedException.ThrowIf(_disposed, this); + if (expectedGeneration is int expected && expected != _runGeneration) + { + return Array.Empty(); + } + + return WriteCurrentRun(snapshot); + } + } + + private IReadOnlyList WriteCurrentRun(SessionSnapshot snapshot) { - ObjectDisposedException.ThrowIf(_disposed, this); Directory.CreateDirectory(_runDirectory); var path = PathFor(snapshot.Id); var tmp = path + ".tmp"; @@ -203,6 +249,112 @@ public void DeleteMany(IEnumerable sessionIds) } } + /// + /// Drops the current run folder from the index and moves it aside. New writes use a fresh folder. + /// The moved folder is deleted on a background thread. Other run folders are left in place. + /// + public void AbandonCurrentRun() + { + string old; + lock (_writeGate) + { + lock (_gate) + { + old = _runDirectory; + Interlocked.Increment(ref _runGeneration); + RemoveIndexEntriesUnder(old); + _runDirectory = CreateRunDirectory(_rootDirectory, DateTimeOffset.UtcNow); + } + + if (!Directory.Exists(old)) + { + return; + } + + var trash = Path.Combine(Path.GetTempPath(), "ti-discard-" + Guid.NewGuid().ToString("N")); + try + { + Directory.Move(old, trash); + QueueCleanup([trash]); + } + catch + { + QueueCleanup([old]); + } + } + } + + /// Removes session files from the size index immediately and deletes them in the background. + public void ScheduleDelete(IEnumerable sessionIds) + { + var paths = new List(); + lock (_writeGate) + { + lock (_gate) + { + foreach (var id in sessionIds) + { + var path = Path.Combine( + _runDirectory, + id.ToString("D", CultureInfo.InvariantCulture) + ".har"); + paths.Add(path); + RemoveIndexEntryLocked(path); + } + } + } + + QueueCleanup(paths); + } + + /// + /// Applies when positive, drops the oldest indexed files from the budget, + /// and deletes those files in the background. Returns current-run session ids that were dropped. + /// + public IReadOnlyList ScheduleBudgetPrune(long maxBytes) + { + List paths; + List currentRunIds; + lock (_writeGate) + { + lock (_gate) + { + if (maxBytes > 0) + { + _maxBytes = maxBytes; + } + + if (_trackedBytes <= _maxBytes) + { + return Array.Empty(); + } + + var ordered = _index + .Select(kv => (Path: kv.Key, kv.Value.SessionId, kv.Value.Length, kv.Value.LastWriteUtc)) + .OrderBy(x => x.LastWriteUtc) + .ToList(); + paths = new List(); + currentRunIds = new List(); + foreach (var entry in ordered) + { + if (_trackedBytes <= _maxBytes) + { + break; + } + + RemoveIndexEntryLocked(entry.Path); + paths.Add(entry.Path); + if (IsUnderRunDirectory(entry.Path)) + { + currentRunIds.Add(entry.SessionId); + } + } + } + } + + QueueCleanup(paths); + return currentRunIds; + } + /// Deletes HARs for this process run only; other run folders stay for Import HAR. public void ClearAll() { @@ -532,6 +684,90 @@ private void RemoveFromIndex(string path, long? trackedLength) } } + /// Caller holds . + private void RemoveIndexEntryLocked(string path) + { + var key = Path.GetFullPath(path); + if (_index.TryGetValue(key, out var prev)) + { + _trackedBytes = Math.Max(0, _trackedBytes - prev.Length); + _index.Remove(key); + } + } + + private void QueueCleanup(IReadOnlyList paths) + { + if (paths.Count == 0) + { + return; + } + + lock (_cleanupGate) + { + foreach (var path in paths) + { + if (!string.IsNullOrEmpty(path)) + { + _cleanupPaths.Enqueue(path); + } + } + + _cleanupTask = _cleanupTask.ContinueWith( + _ => DrainCleanupBatch(), + CancellationToken.None, + TaskContinuationOptions.RunContinuationsAsynchronously, + TaskScheduler.Default); + } + } + + private void DrainCleanupBatch() + { + LastCleanupThreadId = Environment.CurrentManagedThreadId; + List batch; + lock (_cleanupGate) + { + if (_cleanupPaths.Count == 0) + { + return; + } + + batch = _cleanupPaths.ToList(); + _cleanupPaths.Clear(); + } + + foreach (var path in batch) + { + TryDeletePathWithRetry(path); + } + } + + private static void TryDeletePathWithRetry(string path) + { + for (var attempt = 0; attempt < 8; attempt++) + { + try + { + if (Directory.Exists(path)) + { + Directory.Delete(path, recursive: true); + return; + } + + if (File.Exists(path)) + { + File.Delete(path); + return; + } + + return; + } + catch + { + Thread.Sleep(25 * (attempt + 1)); + } + } + } + private void RemoveIndexEntriesUnder(string directory) { var prefix = Path.GetFullPath(directory); diff --git a/src/Titanium.Inspector/Services/SessionListCollection.cs b/src/Titanium.Inspector/Services/SessionListCollection.cs new file mode 100644 index 000000000..ab06b6f0d --- /dev/null +++ b/src/Titanium.Inspector/Services/SessionListCollection.cs @@ -0,0 +1,26 @@ +using System.Collections.ObjectModel; +using System.Collections.Specialized; +using System.ComponentModel; + +namespace Titanium.Inspector.Services; + +/// +/// Session list that can swap its contents with one +/// so the DataGrid does not process one add or remove per row. +/// +public sealed class SessionListCollection : ObservableCollection +{ + public void ReplaceAll(IReadOnlyList items) + { + CheckReentrancy(); + Items.Clear(); + foreach (var item in items) + { + Items.Add(item); + } + + OnPropertyChanged(new PropertyChangedEventArgs(nameof(Count))); + OnPropertyChanged(new PropertyChangedEventArgs("Item[]")); + OnCollectionChanged(new NotifyCollectionChangedEventArgs(NotifyCollectionChangedAction.Reset)); + } +} diff --git a/src/Titanium.Inspector/Services/SessionStore.cs b/src/Titanium.Inspector/Services/SessionStore.cs index 854691ca3..a398905bf 100644 --- a/src/Titanium.Inspector/Services/SessionStore.cs +++ b/src/Titanium.Inspector/Services/SessionStore.cs @@ -14,9 +14,10 @@ public sealed class SessionStore : IDisposable private readonly SessionStoreOptions _options; private readonly SessionBodyDiskCache? _disk; private readonly Dictionary _byId = new(); - private readonly Channel? _spillChannel; + private readonly Channel? _spillChannel; private readonly CancellationTokenSource? _spillCts; private readonly Task? _spillLoop; + private long _spillEpoch; private int _pendingSpills; private long _inMemoryBodyBytes; private int _spilledCount; @@ -33,7 +34,7 @@ public SessionStore(SessionStoreOptions? options = null, string? cacheDirectory var root = cacheDirectory ?? SessionBodyDiskCache.GetDefaultDirectory(); _disk = new SessionBodyDiskCache(root, _options.DiskCacheMaxBytes, TimeSpan.FromDays(7)); // Prior runs stay under other timestamped folders for Import HAR; this run writes here only. - _spillChannel = Channel.CreateUnbounded(new UnboundedChannelOptions + _spillChannel = Channel.CreateUnbounded(new UnboundedChannelOptions { SingleReader = true, SingleWriter = false, @@ -42,7 +43,7 @@ public SessionStore(SessionStoreOptions? options = null, string? cacheDirectory _spillLoop = Task.Run(() => SpillLoopAsync(_spillCts.Token), _spillCts.Token); } - Sessions = new ObservableCollection(); + Sessions = new SessionListCollection(); } public ObservableCollection Sessions { get; } @@ -260,30 +261,80 @@ public void Remove(IEnumerable ids) if (removed.Count > 0) { + _disk?.ScheduleDelete(removed.Select(s => s.Id)); SessionsRemoved?.Invoke(removed); } } + /// + /// Drops every in-memory session and rotates the current-run HAR folder. + /// Does not raise — the caller clears the grid itself. + /// Disk deletion continues in the background (). + /// public void Clear() { ObjectDisposedException.ThrowIf(_disposed, this); - List removed; lock (_gate) { - removed = _byId.Values.ToList(); _byId.Clear(); _inMemoryBodyBytes = 0; _spilledCount = 0; Sessions.Clear(); + Interlocked.Increment(ref _spillEpoch); } - _disk?.ClearAll(); - if (removed.Count > 0) + _disk?.AbandonCurrentRun(); + } + + /// + /// Inserts many sessions under one lock and enforces the memory cap once. + /// Spill writes are queued; HAR files are not written on the caller thread. + /// + public void AddMany(IReadOnlyList snapshots) + { + if (_disposed || snapshots.Count == 0) + { + return; + } + + List? removed = null; + lock (_gate) + { + foreach (var snapshot in snapshots) + { + if (_byId.ContainsKey(snapshot.Id)) + { + MaybeSpillFinishedLocked(snapshot); + continue; + } + + _byId[snapshot.Id] = snapshot; + Sessions.Add(snapshot); + MaybeSpillFinishedLocked(snapshot); + } + + EnforceLimitsLocked(ref removed); + } + + if (removed is { Count: > 0 }) { SessionsRemoved?.Invoke(removed); } } + /// Waits for queued spill writes and background HAR deletes. + public async Task FlushDiskCleanupAsync(TimeSpan? timeout = null) + { + await FlushSpillAsync(timeout).ConfigureAwait(false); + if (_disk is not null) + { + await _disk.FlushCleanupAsync().ConfigureAwait(false); + } + } + + /// Thread that last deleted abandoned HAR files. Zero until a cleanup runs. + internal int LastDiskCleanupThreadId => _disk?.LastCleanupThreadId ?? 0; + public async Task EnsureBodiesLoadedAsync(SessionSnapshot snapshot, CancellationToken ct = default) { ObjectDisposedException.ThrowIf(_disposed, this); @@ -605,6 +656,17 @@ private void MaybeSpillFinishedLocked(SessionSnapshot snapshot) private void EnforceLimitsLocked(ref List? removed) { + var excess = _byId.Count - _options.MaxSessionsInMemory; + if (excess <= 0) + { + return; + } + + if (excess >= 32) + { + EvictOldestBulkLocked(excess, ref removed); + } + while (_byId.Count > _options.MaxSessionsInMemory) { if (!TryEvictOldestLocked(out var evicted)) @@ -617,6 +679,53 @@ private void EnforceLimitsLocked(ref List? removed) } } + private void EvictOldestBulkLocked(int excess, ref List? removed) + { + var evict = new List(excess); + var keep = new List(Math.Max(0, Sessions.Count - excess)); + foreach (var snap in Sessions) + { + var pinned = _pinnedSessionId is long pin && snap.Id == pin; + if (evict.Count < excess && !pinned) + { + evict.Add(snap); + } + else + { + keep.Add(snap); + } + } + + if (evict.Count == 0) + { + return; + } + + foreach (var snap in evict) + { + _byId.Remove(snap.Id); + if (snap.BodiesOnDisk) + { + _spilledCount = Math.Max(0, _spilledCount - 1); + } + + if (_disk is not null && _spillChannel is not null && HasInMemoryBodies(snap)) + { + EnqueueSpillWrite(CloneForDisk(snap)); + } + + removed ??= new List(); + removed.Add(snap); + } + + if (Sessions is SessionListCollection list) + { + list.ReplaceAll(keep); + } + + RecalcInMemoryBodyBytesLocked(); + } + private void QueueSpillLocked(SessionSnapshot snap) { var keepInRam = _pinnedSessionId is long pin && snap.Id == pin; @@ -639,8 +748,10 @@ private void QueueSpillLocked(SessionSnapshot snap) private void EnqueueSpillWrite(SessionSnapshot copy) { + var epoch = Interlocked.Read(ref _spillEpoch); + var generation = _disk?.CurrentGeneration ?? 0; Interlocked.Increment(ref _pendingSpills); - if (!_spillChannel!.Writer.TryWrite(copy)) + if (!_spillChannel!.Writer.TryWrite(new SpillWork(copy, epoch, generation))) { Interlocked.Decrement(ref _pendingSpills); } @@ -737,6 +848,31 @@ private static SessionSnapshot CloneForDisk(SessionSnapshot snap) => private List RemoveIdsLocked(HashSet ids) { var removed = new List(); + if (ids.Count >= 32 && ids.Count * 2 >= Sessions.Count && Sessions is SessionListCollection list) + { + var keep = new List(Math.Max(0, Sessions.Count - ids.Count)); + foreach (var snap in Sessions) + { + if (!ids.Contains(snap.Id)) + { + keep.Add(snap); + continue; + } + + _byId.Remove(snap.Id); + if (snap.BodiesOnDisk) + { + _spilledCount = Math.Max(0, _spilledCount - 1); + } + + removed.Add(snap); + } + + list.ReplaceAll(keep); + RecalcInMemoryBodyBytesLocked(); + return removed; + } + for (var i = Sessions.Count - 1; i >= 0; i--) { var snap = Sessions[i]; @@ -752,7 +888,6 @@ private List RemoveIdsLocked(HashSet ids) _spilledCount = Math.Max(0, _spilledCount - 1); } - _disk?.Delete(snap.Id); removed.Add(snap); } @@ -810,11 +945,16 @@ private async Task SpillLoopAsync(CancellationToken ct) try { - await foreach (var snap in _spillChannel.Reader.ReadAllAsync(ct).ConfigureAwait(false)) + await foreach (var work in _spillChannel.Reader.ReadAllAsync(ct).ConfigureAwait(false)) { try { - var pruned = _disk.Write(snap); + if (work.Epoch != Interlocked.Read(ref _spillEpoch)) + { + continue; + } + + var pruned = _disk.Write(work.Snapshot, work.DiskGeneration); MarkBodiesMissing(pruned); } catch @@ -832,4 +972,6 @@ private async Task SpillLoopAsync(CancellationToken ct) // Shutdown. } } + + private readonly record struct SpillWork(SessionSnapshot Snapshot, long Epoch, int DiskGeneration); } diff --git a/src/Titanium.Inspector/Titanium.Inspector.csproj b/src/Titanium.Inspector/Titanium.Inspector.csproj index 173e859df..49ae8476a 100644 --- a/src/Titanium.Inspector/Titanium.Inspector.csproj +++ b/src/Titanium.Inspector/Titanium.Inspector.csproj @@ -8,7 +8,7 @@ enable true false - 7.0.15 + 7.0.16 Jehonathan Thomas Titanium Inspector desktop traffic debugger (PolyForm Noncommercial). LICENSE diff --git a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs index bfac3295b..84117b783 100644 --- a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs +++ b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs @@ -18,10 +18,23 @@ namespace Titanium.Inspector.ViewModels; public sealed partial class MainWindowViewModel { + private const int BulkGridEditThreshold = 32; + private const string SearchingBodiesStatus = "Searching bodies…"; + + private int _bodyFilterGeneration; + private CancellationTokenSource? _bodyFilterCts; + + /// Managed thread id of the last body: scan. Zero until one runs. + internal int LastBodyFilterThreadId { get; private set; } + + /// Managed thread id of the last HAR/archive export write. + internal int LastExportThreadId { get; private set; } + private Task ClearSessionsAsync() { // Drop selection before mutating the grid so the DataGrid cannot cascade-select // a neighbor row (SelectedSession setter would reopen a closed details pane). + CancelBodyFilter(); _selectedSessions.Clear(); SelectedSession = null; ShowSessionDetails = false; @@ -30,6 +43,7 @@ private Task ClearSessionsAsync() _suppressOpenSessionDetails = true; try { + // Store.Clear does not raise SessionsRemoved. One grid Clear is a single Reset. _store.Clear(); Sessions.Clear(); } @@ -68,14 +82,8 @@ private Task RemoveSelectedSessionsAsync() _suppressOpenSessionDetails = true; try { + // SessionsRemoved updates the grid (one Reset when the selection is large). _store.Remove(ids); - for (var i = Sessions.Count - 1; i >= 0; i--) - { - if (ids.Contains(Sessions[i].Id)) - { - Sessions.RemoveAt(i); - } - } } finally { @@ -582,13 +590,24 @@ private void OnSessionsBatchAdded(IReadOnlyList batch) return; } - var bodyMatcher = BodyMatcherOrNull(); foreach (var snapshot in batch) { _store.Add(snapshot); + } + + if (SessionSearch.HasBodyToken(SearchQuery)) + { + // HAR body reads stay off the UI thread and coalesce while capture is hot. + RunBodyFilter(); + RefreshSessionCountText(); + return; + } + + foreach (var snapshot in batch) + { // SessionUpdated may have raced ahead of this batched capture add and already // inserted the row; never append the same snapshot twice. - if (SessionSearch.Matches(snapshot, SearchQuery, bodyMatcher) && + if (SessionSearch.Matches(snapshot, SearchQuery, bodyMatcher: null) && Sessions.IndexOf(snapshot) < 0) { Sessions.Add(snapshot); @@ -598,15 +617,16 @@ private void OnSessionsBatchAdded(IReadOnlyList batch) RefreshSessionCountText(); } - private Func? BodyMatcherOrNull() => - SessionSearch.HasBodyToken(SearchQuery) - ? (s, needle) => _store.TryMatchBodySearch(s, needle) - : null; - private void OnSessionAddedToFilter(SessionSnapshot snapshot) { + if (SessionSearch.HasBodyToken(SearchQuery)) + { + RunBodyFilter(); + return; + } + // Store already holds the row — append to the filtered grid in place. - if (SessionSearch.Matches(snapshot, SearchQuery, BodyMatcherOrNull()) && + if (SessionSearch.Matches(snapshot, SearchQuery, bodyMatcher: null) && Sessions.IndexOf(snapshot) < 0) { Sessions.Add(snapshot); @@ -629,7 +649,13 @@ private void OnSessionUpdatedForFilter(SessionSnapshot snapshot) return; } - var matches = SessionSearch.Matches(snapshot, SearchQuery, BodyMatcherOrNull()); + if (SessionSearch.HasBodyToken(SearchQuery)) + { + RunBodyFilter(); + return; + } + + var matches = SessionSearch.Matches(snapshot, SearchQuery, bodyMatcher: null); var index = Sessions.IndexOf(snapshot); if (matches) { @@ -671,13 +697,7 @@ private void OnSessionsRemoved(IReadOnlyList removed) _suppressOpenSessionDetails = true; try { - for (var i = Sessions.Count - 1; i >= 0; i--) - { - if (ids.Contains(Sessions[i].Id)) - { - Sessions.RemoveAt(i); - } - } + RemoveVisibleByIds(ids); } finally { @@ -736,17 +756,39 @@ private void NotifyQuickFilterProperties() } private void ApplyFilter() { + if (SessionSearch.HasBodyToken(SearchQuery)) + { + RunBodyFilter(); + return; + } + + CancelBodyFilter(); var previouslySelected = SelectedSession; var detailsWereOpen = ShowSessionDetails; - var matched = SessionSearch.Filter(_all, SearchQuery, BodyMatcherOrNull()).ToList(); + var matched = SessionSearch.Filter(_all, SearchQuery, bodyMatcher: null).ToList(); + ReplaceVisibleSessions(matched); + RestoreSelectionAfterFilter(previouslySelected, detailsWereOpen); + } + + private void ReplaceVisibleSessions(IReadOnlyList matched) + { + if (Sessions is SessionListCollection list) + { + list.ReplaceAll(matched); + return; + } + Sessions.Clear(); foreach (var s in matched) { Sessions.Add(s); } + } + private void RestoreSelectionAfterFilter(SessionSnapshot? previouslySelected, bool detailsWereOpen) + { // Restore single selection used by the detail pane when the row still matches the filter. - // Do not force the pane open — the user may have closed it, and Sessions.Clear() can + // Do not force the pane open — the user may have closed it, and a grid reset can // briefly null SelectedSession via the DataGrid binding. if (previouslySelected is not null && Sessions.Contains(previouslySelected)) { @@ -770,6 +812,127 @@ private void ApplyFilter() SelectedSession = null; } } + + private void RemoveVisibleByIds(HashSet ids) + { + if (ids.Count == 0 || Sessions.Count == 0) + { + return; + } + + var removeCount = 0; + foreach (var session in Sessions) + { + if (ids.Contains(session.Id)) + { + removeCount++; + } + } + + if (removeCount == 0) + { + return; + } + + if (removeCount == Sessions.Count) + { + Sessions.Clear(); + return; + } + + if (removeCount >= BulkGridEditThreshold && removeCount * 2 >= Sessions.Count) + { + var keep = new List(Sessions.Count - removeCount); + foreach (var session in Sessions) + { + if (!ids.Contains(session.Id)) + { + keep.Add(session); + } + } + + ReplaceVisibleSessions(keep); + return; + } + + for (var i = Sessions.Count - 1; i >= 0; i--) + { + if (ids.Contains(Sessions[i].Id)) + { + Sessions.RemoveAt(i); + } + } + } + + private void CancelBodyFilter() + { + _bodyFilterGeneration++; + _bodyFilterCts?.Cancel(); + _bodyFilterCts?.Dispose(); + _bodyFilterCts = null; + } + + private void RunBodyFilter() + { + _bodyFilterCts?.Cancel(); + _bodyFilterCts?.Dispose(); + _bodyFilterCts = new CancellationTokenSource(); + var token = _bodyFilterCts.Token; + var generation = ++_bodyFilterGeneration; + var query = SearchQuery; + var sessions = _all.ToList(); + var previouslySelected = SelectedSession; + var detailsWereOpen = ShowSessionDetails; + + List Match() + { + LastBodyFilterThreadId = Environment.CurrentManagedThreadId; + return SessionSearch.Filter( + sessions, + query, + (snapshot, needle) => _store.TryMatchBodySearch(snapshot, needle)).ToList(); + } + + void Apply(List matched) + { + if (generation != _bodyFilterGeneration || token.IsCancellationRequested) + { + return; + } + + ReplaceVisibleSessions(matched); + RestoreSelectionAfterFilter(previouslySelected, detailsWereOpen); + RefreshSessionCountText(); + if (StatusText == SearchingBodiesStatus) + { + StatusText = StatusReady; + } + } + + // Unit tests have no dispatcher. Match on a dedicated thread (the test thread-pool can + // inline Task.Run onto the caller) and apply before returning. The live UI debounces. + if (Application.Current is null) + { + var matched = RunOnBackgroundThread(Match); + Apply(matched); + return; + } + + StatusText = SearchingBodiesStatus; + _ = Task.Run(async () => + { + try + { + await Task.Delay(200, token).ConfigureAwait(false); + var matched = Match(); + await MarshalToUiAsync(() => Apply(matched), token).ConfigureAwait(false); + } + catch (OperationCanceledException) + { + // A newer query replaced this scan. + } + }, token); + } private async Task ExportHarAsync() { if (_all.Count == 0) @@ -788,18 +951,15 @@ private async Task ExportHarAsync() try { var sessions = _all.ToList(); - // Stay on the UI sync context (RelayCommand). ConfigureAwait(false) + StatusText update - // raced with headless WaitUntil pumps on macOS (file written, StatusText stayed Ready). - SetStatus("Exporting HAR…", StatusSeverity.Busy); - await _store.WithBodiesForExportAsync( + await ExportSessionsOffUiAsync( sessions, - list => SessionArchive.ExportHarAsync(list, path, _statusRevertCts?.Token ?? CancellationToken.None), - _statusRevertCts?.Token ?? CancellationToken.None); - SetOutcomeStatus($"Exported {sessions.Count} sessions to {path}", StatusSeverity.Success, toastImportant: true); + "Exporting HAR…", + list => SessionArchive.ExportHarAsync(list, path, StatusCancelToken), + $"Exported {sessions.Count} sessions to {path}").ConfigureAwait(false); } catch (Exception ex) { - SetOutcomeStatus("Export HAR failed: " + Truncate(ex.Message, 160), StatusSeverity.Error, toastImportant: true); + ReportExportFailure("Export HAR failed: " + Truncate(ex.Message, 160)); } } private async Task ExportSelectedHarAsync() @@ -820,16 +980,15 @@ private async Task ExportSelectedHarAsync() try { - SetStatus("Exporting HAR…", StatusSeverity.Busy); - await _store.WithBodiesForExportAsync( + await ExportSessionsOffUiAsync( sessions, - list => SessionArchive.ExportHarAsync(list, path, _statusRevertCts?.Token ?? CancellationToken.None), - _statusRevertCts?.Token ?? CancellationToken.None); - SetOutcomeStatus($"Exported {sessions.Count} sessions to {path}", StatusSeverity.Success, toastImportant: true); + "Exporting HAR…", + list => SessionArchive.ExportHarAsync(list, path, StatusCancelToken), + $"Exported {sessions.Count} sessions to {path}").ConfigureAwait(false); } catch (Exception ex) { - SetOutcomeStatus("Export HAR failed: " + Truncate(ex.Message, 160), StatusSeverity.Error, toastImportant: true); + ReportExportFailure("Export HAR failed: " + Truncate(ex.Message, 160)); } } private async Task ImportHarAsync() @@ -842,32 +1001,18 @@ private async Task ImportHarAsync() } SetStatus("Importing…", StatusSeverity.Busy); - var imported = new List(); - foreach (var path in paths) - { - if (path.EndsWith(".zip", StringComparison.OrdinalIgnoreCase)) - { - imported.AddRange(await SessionArchive.ImportNativeArchiveAsync( - path, _statusRevertCts?.Token ?? CancellationToken.None)); - } - else - { - imported.AddRange(await SessionArchive.ImportHarAsync( - path, _statusRevertCts?.Token ?? CancellationToken.None)); - } - } - - foreach (var snap in imported) - { - _store.Add(snap); - } - - ApplyFilter(); - RefreshSessionCountText(); + var token = StatusCancelToken; + var imported = await ReadImportedSessionsAsync(paths, token).ConfigureAwait(false); var label = paths.Count == 1 ? Path.GetFileName(paths[0]) : $"{paths.Count} files"; - SetOutcomeStatus($"Appended {imported.Count} sessions from {label}", StatusSeverity.Success, toastImportant: true); + var count = imported.Count; + QueueLiveUi(() => + { + AppendImportedSessions(imported); + RefreshSessionCountText(); + PresentOutcome($"Appended {count} sessions from {label}", StatusSeverity.Success); + }); } private async Task ExportArchiveAsync() { @@ -887,18 +1032,15 @@ private async Task ExportArchiveAsync() try { var sessions = _all.ToList(); - // Stay on the UI sync context (RelayCommand). ConfigureAwait(false) + StatusText update - // raced with headless WaitUntil pumps on macOS (file written, StatusText stayed Ready). - SetStatus("Exporting archive…", StatusSeverity.Busy); - await _store.WithBodiesForExportAsync( + await ExportSessionsOffUiAsync( sessions, - list => SessionArchive.ExportNativeArchiveAsync(list, path, _statusRevertCts?.Token ?? CancellationToken.None), - _statusRevertCts?.Token ?? CancellationToken.None); - SetOutcomeStatus($"Exported {sessions.Count} sessions to {path}", StatusSeverity.Success, toastImportant: true); + "Exporting archive…", + list => SessionArchive.ExportNativeArchiveAsync(list, path, StatusCancelToken), + $"Exported {sessions.Count} sessions to {path}").ConfigureAwait(false); } catch (Exception ex) { - SetOutcomeStatus("Export archive failed: " + Truncate(ex.Message, 160), StatusSeverity.Error, toastImportant: true); + ReportExportFailure("Export archive failed: " + Truncate(ex.Message, 160)); } } private async Task ExportSelectedArchiveAsync() @@ -919,16 +1061,15 @@ private async Task ExportSelectedArchiveAsync() try { - SetStatus("Exporting archive…", StatusSeverity.Busy); - await _store.WithBodiesForExportAsync( + await ExportSessionsOffUiAsync( sessions, - list => SessionArchive.ExportNativeArchiveAsync(list, path, _statusRevertCts?.Token ?? CancellationToken.None), - _statusRevertCts?.Token ?? CancellationToken.None); - SetOutcomeStatus($"Exported {sessions.Count} sessions to {path}", StatusSeverity.Success, toastImportant: true); + "Exporting archive…", + list => SessionArchive.ExportNativeArchiveAsync(list, path, StatusCancelToken), + $"Exported {sessions.Count} sessions to {path}").ConfigureAwait(false); } catch (Exception ex) { - SetOutcomeStatus("Export archive failed: " + Truncate(ex.Message, 160), StatusSeverity.Error, toastImportant: true); + ReportExportFailure("Export archive failed: " + Truncate(ex.Message, 160)); } } private async Task ImportArchiveAsync() @@ -943,24 +1084,230 @@ private async Task ImportArchiveAsync() SetStatus("Importing archive…", StatusSeverity.Busy); try { - // Stay on the UI sync context (RelayCommand). ConfigureAwait(false) + off-thread - // StatusText throws Avalonia "Call from invalid thread" on Windows CI, and - // nested MarshalToUiAsync StatusText updates flaked on macOS headless. - var imported = await SessionArchive.ImportNativeArchiveAsync(path, _statusRevertCts?.Token ?? CancellationToken.None); - foreach (var snap in imported) + var token = StatusCancelToken; + var imported = await ReadImportedSessionsAsync([path], token).ConfigureAwait(false); + var fileName = Path.GetFileName(path); + var count = imported.Count; + QueueLiveUi(() => { - _store.Add(snap); + AppendImportedSessions(imported); + RefreshSessionCountText(); + PresentOutcome($"Appended {count} sessions from {fileName}", StatusSeverity.Success); + }); + } + catch (Exception ex) + { + ReportExportFailure("Import archive failed: " + Truncate(ex.Message, 160)); + } + } + private void ReportExportFailure(string message) => + QueueLiveUi(() => PresentOutcome(message, StatusSeverity.Error)); + + /// + /// Headless Dispatch starts with ResetForUnitTests, which runs queued jobs + /// before IGlobalClock exists. WindowNotificationManager.Show is async void: + /// its card post then throws, and the later Close resumes on the thread pool. + /// Stash the outcome and present it on the next live dispatcher turn. + /// + private int _seenAvaloniaApp; + private int _hasDeferredUi; + private readonly object _deferredUiGate = new(); + private Action? _deferredUi; + + private void NoteAvaloniaApp() + { + if (Application.Current is not null) + { + Volatile.Write(ref _seenAvaloniaApp, 1); + } + } + + private bool RunsWithoutAvaloniaApp => + Application.Current is null && Volatile.Read(ref _seenAvaloniaApp) == 0; + + private bool IsLiveDispatcherTurn() + { + NoteAvaloniaApp(); + // ResetForUnitTests runs while the new locator scope is empty, so Current is null. + // A live Dispatch action runs after SetupUnsafe. + return Application.Current is not null && Dispatcher.UIThread.CheckAccess(); + } + + /// + /// Apply export/import UI updates that were queued off the dispatcher. + /// Headless pumps call this after the application services exist. + /// + public void FlushDeferredInspectorUi() + { + if (!IsLiveDispatcherTurn()) + { + return; + } + + Action? pending; + lock (_deferredUiGate) + { + pending = _deferredUi; + _deferredUi = null; + Volatile.Write(ref _hasDeferredUi, 0); + } + + pending?.Invoke(); + } + + internal void QueueLiveUi(Action action) + { + if (IsLiveDispatcherTurn()) + { + FlushDeferredInspectorUi(); + action(); + return; + } + + if (RunsWithoutAvaloniaApp) + { + action(); + return; + } + + lock (_deferredUiGate) + { + var previous = _deferredUi; + _deferredUi = previous is null + ? action + : () => + { + previous(); + action(); + }; + Volatile.Write(ref _hasDeferredUi, 1); + } + + // Production loop drains this. A headless reset may run it with no clock; + // Flush then returns and the next live turn presents the outcome. + if (Application.Current is not null && !Dispatcher.UIThread.CheckAccess()) + { + Dispatcher.UIThread.Post(FlushDeferredInspectorUi); + } + } + + private void PresentOutcome(string text, StatusSeverity severity) + { + // Show() is async void and awaits a delay before Close. With no sync context that + // Close runs on the thread pool and aborts the headless host. Skip the toast then; + // the status bar text is what the UI tests observe. + var toast = SynchronizationContext.Current is not null; + SetOutcomeStatus(text, severity, toastImportant: toast); + if (!toast || Application.Current is null || !Dispatcher.UIThread.CheckAccess()) + { + return; + } + + // Show() posts the card. Drain it before this dispatch resets the dispatcher, + // while IGlobalClock is still registered. + Dispatcher.UIThread.RunJobs(); + } + + private async Task ExportSessionsOffUiAsync( + IReadOnlyList sessions, + string busyText, + Func, Task> write, + string successText) + { + NoteAvaloniaApp(); + SetStatus(busyText, StatusSeverity.Busy); + var token = StatusCancelToken; + void Write() + { + LastExportThreadId = Environment.CurrentManagedThreadId; + _store.WithBodiesForExportAsync(sessions, write, token).GetAwaiter().GetResult(); + } + + if (RunsWithoutAvaloniaApp) + { + // Unit tests have no dispatcher. Finish before the command returns. + RunOnBackgroundThread(Write); + PresentOutcome(successText, StatusSeverity.Success); + return; + } + + await Task.Run(Write, token).ConfigureAwait(false); + QueueLiveUi(() => PresentOutcome(successText, StatusSeverity.Success)); + } + + private async Task> ReadImportedSessionsAsync( + IReadOnlyList paths, + CancellationToken token) + { + NoteAvaloniaApp(); + List Read() + { + var list = new List(); + foreach (var path in paths) + { + token.ThrowIfCancellationRequested(); + if (path.EndsWith(".zip", StringComparison.OrdinalIgnoreCase)) + { + list.AddRange(SessionArchive.ImportNativeArchiveAsync(path, token).GetAwaiter().GetResult()); + } + else + { + list.AddRange(SessionArchive.ImportHarAsync(path, token).GetAwaiter().GetResult()); + } } - ApplyFilter(); - RefreshSessionCountText(); - SetOutcomeStatus($"Appended {imported.Count} sessions from {Path.GetFileName(path)}", StatusSeverity.Success, toastImportant: true); + return list; } - catch (Exception ex) + + if (RunsWithoutAvaloniaApp) { - SetOutcomeStatus("Import archive failed: " + Truncate(ex.Message, 160), StatusSeverity.Error, toastImportant: true); + return RunOnBackgroundThread>(Read); } + + return await Task.Run(Read, token).ConfigureAwait(false); } + + private static T RunOnBackgroundThread(Func work) + { + T? result = default; + Exception? error = null; + var thread = new Thread(() => + { + try + { + result = work(); + } + catch (Exception ex) + { + error = ex; + } + }) + { + IsBackground = true, + }; + thread.Start(); + thread.Join(); + if (error is not null) + { + throw error; + } + + return result!; + } + + private static void RunOnBackgroundThread(Action work) => + RunOnBackgroundThread(() => + { + work(); + return true; + }); + + private void AppendImportedSessions(List imported) + { + _store.AddMany(imported); + ApplyFilter(); + } + private IReadOnlyList ResolveExportSelection() { if (_selectedSessions.Count > 0) diff --git a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Updates.cs b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Updates.cs index 547b4156e..87295ab31 100644 --- a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Updates.cs +++ b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Updates.cs @@ -190,7 +190,7 @@ private async Task SendComposerAsync() }; _store.Add(snap); - ApplyFilter(); + OnSessionAddedToFilter(snap); RefreshSessionCountText(); // Select the synthetic row without forcing Inspect open (Composer may already be showing). SelectSessionWithoutOpeningDetails(snap); diff --git a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.cs b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.cs index 9856ea405..339a3530f 100644 --- a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.cs +++ b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.cs @@ -149,6 +149,7 @@ public MainWindowViewModel( public MainWindowViewModel(InspectorViewModelServices services) { + NoteAvaloniaApp(); _buffer = services.Buffer; _registry = services.Registry; _store = services.Registry.Store; @@ -159,7 +160,7 @@ public MainWindowViewModel(InspectorViewModelServices services) _dialogs = services.Dialogs ?? new AvaloniaInspectorDialogs(); _pathPicker = services.PathPicker ?? new AvaloniaInspectorPathPicker(); _statusNotifier = services.StatusNotifier ?? NullStatusNotifier.Instance; - Sessions = new ObservableCollection(); + Sessions = new SessionListCollection(); Breakpoints = new BreakpointViewModel(); AutoResponder = new AutoResponderViewModel(); MapRemote = new MapRemoteViewModel(); @@ -2496,7 +2497,17 @@ public int SelectedDetailTabIndex public string StatusText { - get => _statusText; + get + { + // Headless polls read this after SetupUnsafe. Present any export/import + // result that was stashed while ResetForUnitTests was tearing the clock down. + if (Volatile.Read(ref _hasDeferredUi) != 0) + { + FlushDeferredInspectorUi(); + } + + return _statusText; + } set { if (_settingStatus) diff --git a/src/Titanium.Plus/Titanium.Plus.csproj b/src/Titanium.Plus/Titanium.Plus.csproj index 32608ede2..c854d11ae 100644 --- a/src/Titanium.Plus/Titanium.Plus.csproj +++ b/src/Titanium.Plus/Titanium.Plus.csproj @@ -7,7 +7,7 @@ enable True StrongNameKey.snk - 7.0.15 + 7.0.16 Jehonathan Thomas Titanium Web Proxy Plus advanced features plugin (PolyForm Noncommercial). LICENSE diff --git a/src/Titanium.Web.Proxy.Abstractions/Titanium.Web.Proxy.Abstractions.csproj b/src/Titanium.Web.Proxy.Abstractions/Titanium.Web.Proxy.Abstractions.csproj index af2705780..baa60f1e0 100644 --- a/src/Titanium.Web.Proxy.Abstractions/Titanium.Web.Proxy.Abstractions.csproj +++ b/src/Titanium.Web.Proxy.Abstractions/Titanium.Web.Proxy.Abstractions.csproj @@ -7,7 +7,7 @@ enable True StrongNameKey.snk - 7.0.15 + 7.0.16 Jehonathan Thomas Shared contracts for Titanium Web Proxy routing, clusters, middleware, and plugins. MIT diff --git a/src/Titanium.Web.Proxy.Configuration/Titanium.Web.Proxy.Configuration.csproj b/src/Titanium.Web.Proxy.Configuration/Titanium.Web.Proxy.Configuration.csproj index c3a4b96eb..1718c2646 100644 --- a/src/Titanium.Web.Proxy.Configuration/Titanium.Web.Proxy.Configuration.csproj +++ b/src/Titanium.Web.Proxy.Configuration/Titanium.Web.Proxy.Configuration.csproj @@ -7,7 +7,7 @@ enable True StrongNameKey.snk - 7.0.15 + 7.0.16 Jehonathan Thomas YAML/JSON configuration binding for Titanium Web Proxy CLI and reverse-proxy documents. MIT diff --git a/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs b/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs index 9dbaf57d0..1e76408fb 100644 --- a/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs +++ b/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs @@ -42,14 +42,20 @@ public partial class ProxyServer /// Await-safe shell pool for terminate-lite. ThreadStatic is unsafe here: a second Rent on the /// same worker can while the first request /// still awaits across the shared Response (HTTP/0.0 0 / spliced status lines under load). + /// A shell is owned by exactly one request between Rent (dequeued) and Release (enqueued). + /// is lock-free; ConcurrentBag.Count freezes every + /// per-thread list under Monitor and Rent/Release usually run on different workers, so the bag + /// steal path serialized the hot path (about 7% of proxy-tree CPU in a Linux profile). /// - private static readonly ConcurrentBag H1TerminateLiteClients = new(); - private const int H1TerminateLiteClientPoolCap = 256; + private static readonly ConcurrentQueue H1TerminateLiteClients = new(); + private static int h1TerminateLiteClientCount; + internal const int H1TerminateLiteClientPoolCap = 256; - private static HttpWebClient RentH1TerminateLiteClient(Request request) + internal static HttpWebClient RentH1TerminateLiteClient(Request request) { - if (H1TerminateLiteClients.TryTake(out var client)) + if (H1TerminateLiteClients.TryDequeue(out var client)) { + Interlocked.Decrement(ref h1TerminateLiteClientCount); client.RebindForTerminateLite(request); return client; } @@ -60,12 +66,14 @@ private static HttpWebClient RentH1TerminateLiteClient(Request request) return client; } - private static void ReleaseH1TerminateLiteClient(HttpWebClient client) + internal static void ReleaseH1TerminateLiteClient(HttpWebClient client) { // Drop the origin socket reference so TcpConnectionFactory.Release remains the sole owner. client.RebindForTerminateLite(client.Request); - if (H1TerminateLiteClients.Count < H1TerminateLiteClientPoolCap) - H1TerminateLiteClients.Add(client); + if (Interlocked.Increment(ref h1TerminateLiteClientCount) <= H1TerminateLiteClientPoolCap) + H1TerminateLiteClients.Enqueue(client); + else + Interlocked.Decrement(ref h1TerminateLiteClientCount); } /// diff --git a/src/Titanium.Web.Proxy/Http2/Http2FlowController.cs b/src/Titanium.Web.Proxy/Http2/Http2FlowController.cs index f37ed755d..9275621ce 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2FlowController.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2FlowController.cs @@ -40,6 +40,14 @@ internal sealed class Http2FlowController /// RFC 7540 §6.9.2 default initial flow-control window size for both the connection and every stream. internal const int InitialConnectionWindow = 65535; + /// + /// Smallest DATA payload this proxy will emit just to fit a short send window. Below this, the + /// caller waits for a full frame instead of spraying tiny frames. 4 KiB is large enough to avoid + /// a frame storm and small enough to use the 16,383 bytes left after three 16 KiB frames in a + /// 65,535-byte window. + /// + internal const int PartialDataFrameFloor = 4096; + /// RFC 7540 §6.9.1 - a flow-control window (connection or stream) must never exceed this value. internal const long MaxWindow = int.MaxValue; // 2^31 - 1 @@ -140,6 +148,47 @@ public bool OnWindowUpdate(int streamId, int increment) } } + /// + /// Read-only snapshot of send credit: the minimum of the connection window and the stream window, + /// clamped at zero. Does not reserve. An unknown stream returns 0. The value can be stale by the + /// time the caller reserves — / stay the arbiter. + /// + public int AvailableSendCredit(int streamId) + { + lock (gate) + { + if (!streamWindows.TryGetValue(streamId, out var streamWindow)) + return 0; + + var available = Math.Min(connectionWindow, streamWindow); + if (available <= 0) + return 0; + if (available > int.MaxValue) + return int.MaxValue; + return (int)available; + } + } + + /// + /// How many payload bytes to read for the next DATA frame. A full frame when credit covers it, + /// or when credit is below (the caller then waits in + /// ). When credit is at least that floor and short of a full frame, + /// returns the credit so the frame is not parked for a few bytes. DATA shorter than + /// MAX_FRAME_SIZE is valid (RFC 9113 §4.1). + /// + internal static int SelectDataPayloadCap(int maxFrameSize, long remaining, int availableCredit) + { + if (maxFrameSize <= 0) + maxFrameSize = 16384; + if (remaining <= 0) + return 0; + + var cap = remaining >= maxFrameSize ? maxFrameSize : (int)remaining; + if (availableCredit >= PartialDataFrameFloor && availableCredit < cap) + return availableCredit; + return cap; + } + /// /// Non-blocking variant of . Returns false when credit is /// insufficient instead of waiting. Used to keep HEADERS + first DATA under one write lock diff --git a/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs b/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs index eb0948a6c..29f36e403 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs @@ -1326,7 +1326,13 @@ public override async ValueTask WriteAsync(ReadOnlyMemory buffer, if (expectedLength >= 0 && remaining <= 0) break; - var payloadCap = (int)Math.Min(maxFrameSize, remaining); + // Size the frame to currently advertised credit when that credit is a usable + // partial frame (for example 16,383 bytes left in a 65,535 window). ReserveAsync + // is still all-or-nothing and remains the accounting arbiter if the snapshot is stale. + var payloadCap = Http2FlowController.SelectDataPayloadCap( + maxFrameSize, remaining, flow.AvailableSendCredit(streamId)); + if (payloadCap <= 0) + break; var rented = ArrayPool.Shared.Rent(9 + payloadCap); var read = 0; try diff --git a/src/Titanium.Web.Proxy/Http3/Http3Frame.cs b/src/Titanium.Web.Proxy/Http3/Http3Frame.cs index d7c926ff9..cd3979af3 100644 --- a/src/Titanium.Web.Proxy/Http3/Http3Frame.cs +++ b/src/Titanium.Web.Proxy/Http3/Http3Frame.cs @@ -1,5 +1,6 @@ using System; using System.Buffers; +using System.Collections.Concurrent; using System.IO; using System.Net.Quic; using System.Threading; @@ -97,10 +98,18 @@ public static ValueTask WriteAsync( ulong frameType, ReadOnlyMemory payload, CancellationToken cancellationToken, - bool completeWrites = false) + bool completeWrites = false, + Http3FrameScratch? scratch = null) { // Max VarInt is 8 bytes each for type + length. const int headerCap = 16; + if (scratch != null + && payload.Length <= Http3FrameScratch.Capacity - headerCap + && scratch.TryAcquire()) + { + return WriteScratchAsync(stream, frameType, payload, completeWrites, scratch, cancellationToken); + } + if (payload.Length <= 256) { var rented = ArrayPool.Shared.Rent(headerCap + payload.Length); @@ -123,6 +132,78 @@ public static ValueTask WriteAsync( return WriteLargeAsync(stream, frameType, payload, completeWrites, cancellationToken); } + /// + /// Copies the frame into and does not release it until the + /// write has been consumed. Sync completion calls GetResult + /// before release so MsQuic cannot still hold the memory (the e781b009 ArrayPool bug). + /// A second write while this one is in flight must not call + /// successfully; callers fall back to . + /// + private static ValueTask WriteScratchAsync( + Stream stream, + ulong frameType, + ReadOnlyMemory payload, + bool completeWrites, + Http3FrameScratch scratch, + CancellationToken cancellationToken) + { + int total; + try + { + var span = scratch.Buffer.AsSpan(); + var typeLen = Http3VarInt.Write(span, frameType); + var lengthLen = Http3VarInt.Write(span.Slice(typeLen), (ulong)payload.Length); + var headerLen = typeLen + lengthLen; + if (!payload.IsEmpty) + payload.Span.CopyTo(span.Slice(headerLen)); + total = headerLen + payload.Length; + } + catch + { + scratch.Release(); + throw; + } + + ValueTask write; + try + { + write = WriteBufferAsync(stream, scratch.Buffer.AsMemory(0, total), completeWrites, cancellationToken); + } + catch + { + scratch.Release(); + throw; + } + + if (write.IsCompletedSuccessfully) + { + try + { + write.GetAwaiter().GetResult(); + } + finally + { + scratch.Release(); + } + + return default; + } + + return AwaitScratchAsync(write, scratch); + } + + private static async ValueTask AwaitScratchAsync(ValueTask write, Http3FrameScratch scratch) + { + try + { + await write.ConfigureAwait(false); + } + finally + { + scratch.Release(); + } + } + private static async ValueTask WriteLargeAsync( Stream stream, ulong frameType, @@ -224,3 +305,49 @@ private static ValueTask WriteBufferAsync( } #pragma warning restore CA1416 } + +/// +/// Single-owner buffer for small HTTP/3 frames on one stream. The buffer stays owned until +/// runs, which is only after the write ValueTask has been consumed. +/// Do not return this object to the pool while is held. +/// Writes on one stream are sequential (each WriteAsync is awaited). A concurrent +/// writer fails and uses instead. +/// +internal sealed class Http3FrameScratch +{ + internal const int Capacity = 1024; + + private static readonly ConcurrentBag Pool = new(); + + private readonly byte[] buffer = new byte[Capacity]; + private int busy; + private int checkedOut; + + internal byte[] Buffer => buffer; + + internal static Http3FrameScratch Rent() + { + if (!Pool.TryTake(out var scratch)) + scratch = new Http3FrameScratch(); + scratch.checkedOut = 1; + scratch.busy = 0; + return scratch; + } + + /// + /// Returns the scratch to the pool. If a write is still in flight the instance is dropped + /// instead of being reused, so MsQuic cannot observe a recycled buffer. + /// + internal void Return() + { + if (Volatile.Read(ref busy) != 0) + return; + if (Interlocked.Exchange(ref checkedOut, 0) != 1) + return; + Pool.Add(this); + } + + internal bool TryAcquire() => Interlocked.CompareExchange(ref busy, 1, 0) == 0; + + internal void Release() => Volatile.Write(ref busy, 0); +} diff --git a/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs b/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs index 776701acb..cce2cf555 100644 --- a/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs +++ b/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs @@ -1168,8 +1168,17 @@ private static async Task SendSimpleStatusResponseAsync( { var headers = new List<(string, string)> { (":status", statusCode.ToString()) }; var encoded = QpackEncoder.Encode(headers, qpackContext); - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, encoded, ct, completeWrites: true); - await stream.FlushAsync(ct); + var scratch = Http3FrameScratch.Rent(); + try + { + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, encoded, ct, completeWrites: true, + scratch: scratch); + await stream.FlushAsync(ct); + } + finally + { + scratch.Return(); + } // completeWrites:true already FINed the QuicStream write side. } @@ -1191,8 +1200,16 @@ private static async Task SendInterimResponseAsync( headers.Add((name, header.Value)); } var encoded = QpackEncoder.Encode(headers, qpackContext); - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, encoded, ct); - await stream.FlushAsync(ct); + var scratch = Http3FrameScratch.Rent(); + try + { + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, encoded, ct, scratch: scratch); + await stream.FlushAsync(ct); + } + finally + { + scratch.Return(); + } } /// @@ -1205,33 +1222,44 @@ private static async ValueTask SendPreencodedResponseAsync( Func? streamBodyWriter, CancellationToken ct) { - if (streamBodyWriter != null) + var scratch = Http3FrameScratch.Rent(); + try { - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct); - var bodyWriter = new Http3DataBodyWriter(stream); - await streamBodyWriter(bodyWriter, ct); - await stream.FlushAsync(ct); - return; - } + if (streamBodyWriter != null) + { + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, scratch: scratch); + var bodyWriter = new Http3DataBodyWriter(stream, scratch); + await streamBodyWriter(bodyWriter, ct); + await stream.FlushAsync(ct); + return; + } - if (body.Length >= 16 * 1024) - { - await Http3Frame.WriteHeadersAndDataAsync(stream, qpackHeaders, body, ct, completeWrites: true); - await stream.FlushAsync(ct); - return; - } + if (body.Length >= 16 * 1024) + { + // Too big for the scratch. Keep the pooled coalesce path. + await Http3Frame.WriteHeadersAndDataAsync(stream, qpackHeaders, body, ct, completeWrites: true); + await stream.FlushAsync(ct); + return; + } - if (body.Length > 0) - { - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct); - await Http3Frame.WriteAsync(stream, Http3FrameType.Data, body, ct, completeWrites: true); + if (body.Length > 0) + { + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, scratch: scratch); + await Http3Frame.WriteAsync(stream, Http3FrameType.Data, body, ct, completeWrites: true, + scratch: scratch); + } + else + { + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, completeWrites: true, + scratch: scratch); + } + + await stream.FlushAsync(ct); } - else + finally { - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, completeWrites: true); + scratch.Return(); } - - await stream.FlushAsync(ct); } /// @@ -1253,6 +1281,22 @@ private static async Task SendResponseAsync(QuicStream stream, Response response return; } + var scratch = Http3FrameScratch.Rent(); + try + { + await SendBufferedResponseAsync(stream, response, qpackHeaders, qpackContext, hasTrailers, ct, scratch); + } + finally + { + scratch.Return(); + } + } + + private static async Task SendBufferedResponseAsync( + QuicStream stream, Response response, byte[] qpackHeaders, QpackContext? qpackContext, bool hasTrailers, + CancellationToken ct, Http3FrameScratch scratch) + { + // Send body if present. Ok()/Respond assign Body without setting IsBodyRead (H1 uses // BodyAvailable); requiring IsBodyRead alone dropped every synthetic H3 response body. // Match H1 WriteResponseAsync / H2 EmitBuffered: recompress when Body is plain and @@ -1282,19 +1326,22 @@ private static async Task SendResponseAsync(QuicStream stream, Response response if (body is { Length: > 0 }) { - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct); - await Http3Frame.WriteAsync(stream, Http3FrameType.Data, body, ct, completeWrites: !hasTrailers); + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, scratch: scratch); + await Http3Frame.WriteAsync(stream, Http3FrameType.Data, body, ct, completeWrites: !hasTrailers, + scratch: scratch); } else { - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, completeWrites: !hasTrailers); + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, completeWrites: !hasTrailers, + scratch: scratch); } if (hasTrailers) { var trailerBlock = QpackEncoder.Encode( response.TrailingHeaders.Select(h => (h.Name, h.Value)), qpackContext); - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, trailerBlock, ct, completeWrites: true); + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, trailerBlock, ct, completeWrites: true, + scratch: scratch); } await stream.FlushAsync(ct); @@ -1307,18 +1354,27 @@ private static async Task SendStreamedResponseAsync( QuicStream stream, Response response, ReadOnlyMemory qpackHeaders, QpackContext? qpackContext, bool hasTrailers, CancellationToken ct) { - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct); - // Http3OriginBridge streams the origin body; drain it as DATA frames (same contract as - // H1 BodyStreamWriter / H2 EmitSyntheticResponseAsync). - // Http3DataBodyWriter never FINs (completeWrites stays false) so trailers can follow. - var bodyWriter = new Http3DataBodyWriter(stream); - await response.StreamBodyWriter!(bodyWriter, ct); - response.IsBodySent = true; - if (hasTrailers) + var scratch = Http3FrameScratch.Rent(); + try { - var trailerBlock = QpackEncoder.Encode( - response.TrailingHeaders.Select(h => (h.Name, h.Value)), qpackContext); - await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, trailerBlock, ct, completeWrites: true); + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, qpackHeaders, ct, scratch: scratch); + // Http3OriginBridge streams the origin body; drain it as DATA frames (same contract as + // H1 BodyStreamWriter / H2 EmitSyntheticResponseAsync). + // Http3DataBodyWriter never FINs (completeWrites stays false) so trailers can follow. + var bodyWriter = new Http3DataBodyWriter(stream, scratch); + await response.StreamBodyWriter!(bodyWriter, ct); + response.IsBodySent = true; + if (hasTrailers) + { + var trailerBlock = QpackEncoder.Encode( + response.TrailingHeaders.Select(h => (h.Name, h.Value)), qpackContext); + await Http3Frame.WriteAsync(stream, Http3FrameType.Headers, trailerBlock, ct, completeWrites: true, + scratch: scratch); + } + } + finally + { + scratch.Return(); } // Always Flush — Darwin MsQuic requires an explicit Flush before FIN is observed. @@ -1357,8 +1413,13 @@ private static bool HasUpperAscii(string s) private sealed class Http3DataBodyWriter : Stream { private readonly QuicStream _stream; + private readonly Http3FrameScratch? _scratch; - public Http3DataBodyWriter(QuicStream stream) => _stream = stream; + public Http3DataBodyWriter(QuicStream stream, Http3FrameScratch? scratch = null) + { + _stream = stream; + _scratch = scratch; + } public override bool CanRead => false; public override bool CanSeek => false; @@ -1394,7 +1455,7 @@ public override ValueTask WriteAsync(ReadOnlyMemory buffer, CancellationToken cancellationToken = default) { if (buffer.IsEmpty) return default; - return Http3Frame.WriteAsync(_stream, Http3FrameType.Data, buffer, cancellationToken); + return Http3Frame.WriteAsync(_stream, Http3FrameType.Data, buffer, cancellationToken, scratch: _scratch); } } diff --git a/src/Titanium.Web.Proxy/Properties/AssemblyInfo.cs b/src/Titanium.Web.Proxy/Properties/AssemblyInfo.cs index 0a2427c0c..f834194ac 100644 --- a/src/Titanium.Web.Proxy/Properties/AssemblyInfo.cs +++ b/src/Titanium.Web.Proxy/Properties/AssemblyInfo.cs @@ -77,5 +77,5 @@ // file-properties version disagreed with the package it was published in. Keep both of the values // below equal to (as Major.Minor.Build.0) whenever that property changes. -[assembly: AssemblyVersion("7.0.15.0")] -[assembly: AssemblyFileVersion("7.0.15.0")] +[assembly: AssemblyVersion("7.0.16.0")] +[assembly: AssemblyFileVersion("7.0.16.0")] diff --git a/src/Titanium.Web.Proxy/Titanium.Web.Proxy.csproj b/src/Titanium.Web.Proxy/Titanium.Web.Proxy.csproj index 93c989b46..6eb89e241 100644 --- a/src/Titanium.Web.Proxy/Titanium.Web.Proxy.csproj +++ b/src/Titanium.Web.Proxy/Titanium.Web.Proxy.csproj @@ -13,7 +13,7 @@ - 7.0.15 + 7.0.16