diff --git a/.config/nextest.toml b/.config/nextest.toml new file mode 100644 index 000000000..0abc58828 --- /dev/null +++ b/.config/nextest.toml @@ -0,0 +1,8 @@ +# GitHub-hosted ubuntu runner 当前按 4 vCPU 配置;固定线程数可避免 runner +# 规格变化时测试并发和内存峰值随之漂移。 +[profile.default] +test-threads = 4 + +# 60 秒后标记慢测试,连续两轮仍未结束则终止,避免单个卡死用例拖满整个 job。 +# 真实连接/数据库测试仍有足够时间完成;超时结果保持失败,不隐藏回归。 +slow-timeout = { period = "60s", terminate-after = 2, grace-period = "10s" } diff --git a/.github/workflows/rust-ci.yml b/.github/workflows/rust-ci.yml index a01af779f..c8d6dfdc2 100644 --- a/.github/workflows/rust-ci.yml +++ b/.github/workflows/rust-ci.yml @@ -6,61 +6,7 @@ on: branches: - master - main - paths: - - "Cargo.toml" - - "Cargo.lock" - - "crates/**" - - "apps/**" - - "install.sh" - - "deploy.sh" - - "update.sh" - - "generate_keys.sh" - - ".env.example" - - "README.md" - - "Dockerfile.app" - - "docker-compose.yml" - - "docker-compose.single-node.yml" - - "docker-compose.local.yml" - - "docker-compose.release-local.yml" - - "tests/compose_database_config_test.py" - - "tests/install_*_test.sh" - - "tests/deploy_*_test.sh" - - "tests/update_*_test.sh" - - "tests/release_supply_chain_test.sh" - - "tests/tunnel_installer_config_security_test.sh" - - ".github/workflows/build-tunnel.yml" - - ".github/workflows/deploy-pages.yml" - - ".github/workflows/release.yml" - - ".github/workflows/rust-ci.yml" - - ".github/workflows/nightly.yml" pull_request: - paths: - - "Cargo.toml" - - "Cargo.lock" - - "crates/**" - - "apps/**" - - "install.sh" - - "deploy.sh" - - "update.sh" - - "generate_keys.sh" - - ".env.example" - - "README.md" - - "Dockerfile.app" - - "docker-compose.yml" - - "docker-compose.single-node.yml" - - "docker-compose.local.yml" - - "docker-compose.release-local.yml" - - "tests/compose_database_config_test.py" - - "tests/install_*_test.sh" - - "tests/deploy_*_test.sh" - - "tests/update_*_test.sh" - - "tests/release_supply_chain_test.sh" - - "tests/tunnel_installer_config_security_test.sh" - - ".github/workflows/build-tunnel.yml" - - ".github/workflows/deploy-pages.yml" - - ".github/workflows/release.yml" - - ".github/workflows/rust-ci.yml" - - ".github/workflows/nightly.yml" concurrency: group: rust-ci-${{ github.event_name }}-${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }} @@ -76,8 +22,58 @@ env: CARGO_TERM_COLOR: always jobs: + changes: + name: Detect Rust CI scope + runs-on: ubuntu-latest + outputs: + rust: ${{ steps.scope.outputs.rust }} + shell: ${{ steps.scope.outputs.shell }} + steps: + - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 + with: + fetch-depth: 0 + + - name: Classify changed paths + id: scope + shell: bash + run: | + # workflow_call(Nightly)仍须完整执行;普通 push/PR 只按源码和构建 + # 指纹触发 Rust jobs,安装脚本、Compose、README 等由 shell scope 覆盖。 + if [ "$GITHUB_EVENT_NAME" = "workflow_call" ]; then + echo "rust=true" >> "$GITHUB_OUTPUT" + echo "shell=true" >> "$GITHUB_OUTPUT" + exit 0 + fi + + if [ "$GITHUB_EVENT_NAME" = "pull_request" ]; then + git fetch --no-tags origin "$GITHUB_BASE_REF" --depth=1 + changed_paths=$(git diff --name-only "origin/$GITHUB_BASE_REF...$GITHUB_SHA") + elif [ "$GITHUB_EVENT_NAME" = "push" ] && [ "$GITHUB_EVENT_BEFORE" != "0000000000000000000000000000000000000000" ]; then + changed_paths=$(git diff --name-only "$GITHUB_EVENT_BEFORE" "$GITHUB_SHA") + else + changed_paths=$(git ls-files) + fi + + rust=false + shell=false + while IFS= read -r path; do + case "$path" in + Cargo.toml|Cargo.lock|rust-toolchain.toml|.cargo/*|*.rs|*/Cargo.toml|*/build.rs|*.sql|.github/workflows/*.yml|.github/workflows/*.yaml) + rust=true + ;; + *.sh|*.py|README.md|*/README.md|.env.example|Dockerfile*|docker-compose*.yml|docker-compose*.yaml) + shell=true + ;; + esac + done <<< "$changed_paths" + + echo "rust=$rust" >> "$GITHUB_OUTPUT" + echo "shell=$shell" >> "$GITHUB_OUTPUT" + shell_security: name: Shell security fixtures + needs: changes + if: ${{ needs.changes.outputs.shell == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 @@ -99,6 +95,8 @@ jobs: fmt: name: Format + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 @@ -114,6 +112,8 @@ jobs: clippy_gateway: name: Clippy (Gateway) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 @@ -127,7 +127,9 @@ jobs: - name: Rust cache uses: Swatinem/rust-cache@49a0bdc70d2e1b713ca9e2869b211fcce03d3c1c # v2 with: - shared-key: rust-ci-${{ runner.os }} + # Gateway lint 与 Gateway 测试都可能触发 mold/大型链接依赖,单独隔离缓存 + # 指纹,避免不同 job 的构建产物互相驱逐或复用错误的链接参数。 + shared-key: rust-ci-gateway-clippy-${{ runner.os }} workspaces: . -> target - name: Setup sccache @@ -148,6 +150,8 @@ jobs: clippy_data: name: Clippy (Data) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 @@ -182,6 +186,8 @@ jobs: clippy_rest: name: Clippy (Workspace Rest) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 @@ -218,6 +224,7 @@ jobs: name: Clippy runs-on: ubuntu-latest needs: + - changes - clippy_gateway - clippy_data - clippy_rest @@ -225,6 +232,10 @@ jobs: steps: - name: Verify clippy jobs run: | + if [ "${{ needs.changes.outputs.rust }}" != "true" ]; then + echo "Rust scope unchanged; clippy jobs skipped" + exit 0 + fi if [ "${{ needs.clippy_gateway.result }}" != "success" ] || \ [ "${{ needs.clippy_data.result }}" != "success" ] || \ [ "${{ needs.clippy_rest.result }}" != "success" ]; then @@ -234,12 +245,24 @@ jobs: test_gateway: name: Test (Gateway) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest + # 构建指纹提到 job 级:mold RUSTFLAGS / 栈 / sccache 对 lib、bins、integration 三步保持一致, + # 避免 step 级 env 漂移导致同 job 内 rustc 指纹不一致。 + env: + RUSTC_WRAPPER: sccache + SCCACHE_GHA_ENABLED: "true" + RUST_MIN_STACK: "16777216" + RUSTFLAGS: "-C link-arg=-fuse-ld=mold" steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 - name: Install Rust toolchain uses: dtolnay/rust-toolchain@4360b52568e2003a75bf9bc1d59f33a8e3fc893c # stable + with: + # 与 rust-toolchain.toml、fmt/clippy 钉在同一版本,避免浮动 stable 换指纹导致全量重编 + toolchain: 1.95.0 - name: Show Rust toolchain run: rustup show active-toolchain @@ -247,7 +270,8 @@ jobs: - name: Rust cache uses: Swatinem/rust-cache@49a0bdc70d2e1b713ca9e2869b211fcce03d3c1c # v2 with: - shared-key: rust-ci-${{ runner.os }} + # mold RUSTFLAGS 只在本 job 生效:独立 cache key,避免与无 mold 的 job 互相污染指纹 + shared-key: rust-ci-gateway-test-${{ runner.os }} workspaces: . -> target - name: Setup sccache @@ -263,36 +287,34 @@ jobs: run: pg_config --bindir >> "$GITHUB_PATH" - name: Test lib - env: - RUSTC_WRAPPER: sccache - SCCACHE_GHA_ENABLED: "true" - RUST_MIN_STACK: "16777216" - RUSTFLAGS: "-C link-arg=-fuse-ld=mold" run: cargo nextest run -p aether-gateway --lib - name: Test bins - env: - RUSTC_WRAPPER: sccache - SCCACHE_GHA_ENABLED: "true" - RUST_MIN_STACK: "16777216" - RUSTFLAGS: "-C link-arg=-fuse-ld=mold" run: cargo nextest run -p aether-gateway --bins + # 只运行独立 integration targets;显式列出目标,避免 --tests 再次执行 lib/bin 测试。 + - name: Test integration targets + run: >- + cargo nextest run -p aether-gateway + --test admin_unsigned_identity_headers + --test architecture_guard + - name: Show sccache stats if: always() - env: - RUSTC_WRAPPER: sccache - SCCACHE_GHA_ENABLED: "true" run: sccache --show-stats test_data: name: Test (Data) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 - name: Install Rust toolchain uses: dtolnay/rust-toolchain@4360b52568e2003a75bf9bc1d59f33a8e3fc893c # stable + with: + toolchain: 1.95.0 - name: Show Rust toolchain run: rustup show active-toolchain @@ -328,6 +350,8 @@ jobs: check_data_features: name: Check (Data Feature - ${{ matrix.feature }}) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest strategy: fail-fast: false @@ -340,6 +364,8 @@ jobs: - name: Install Rust toolchain uses: dtolnay/rust-toolchain@4360b52568e2003a75bf9bc1d59f33a8e3fc893c # stable + with: + toolchain: 1.95.0 - name: Rust cache uses: Swatinem/rust-cache@49a0bdc70d2e1b713ca9e2869b211fcce03d3c1c # v2 @@ -365,12 +391,16 @@ jobs: test_rest: name: Test (Workspace Rest) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 - name: Install Rust toolchain uses: dtolnay/rust-toolchain@4360b52568e2003a75bf9bc1d59f33a8e3fc893c # stable + with: + toolchain: 1.95.0 - name: Show Rust toolchain run: rustup show active-toolchain @@ -402,6 +432,8 @@ jobs: test_data_adapters: name: Test (Data Adapter - ${{ matrix.package }}) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest strategy: fail-fast: false @@ -413,6 +445,8 @@ jobs: - name: Install Rust toolchain uses: dtolnay/rust-toolchain@4360b52568e2003a75bf9bc1d59f33a8e3fc893c # stable + with: + toolchain: 1.95.0 - name: Rust cache uses: Swatinem/rust-cache@49a0bdc70d2e1b713ca9e2869b211fcce03d3c1c # v2 @@ -441,12 +475,16 @@ jobs: check_integration_scenarios: name: Test (Integration Scenarios) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 - name: Install Rust toolchain uses: dtolnay/rust-toolchain@4360b52568e2003a75bf9bc1d59f33a8e3fc893c # stable + with: + toolchain: 1.95.0 - name: Rust cache uses: Swatinem/rust-cache@49a0bdc70d2e1b713ca9e2869b211fcce03d3c1c # v2 @@ -477,6 +515,7 @@ jobs: name: Test runs-on: ubuntu-latest needs: + - changes - test_gateway - test_data - check_data_features @@ -487,6 +526,10 @@ jobs: steps: - name: Verify test jobs run: | + if [ "${{ needs.changes.outputs.rust }}" != "true" ]; then + echo "Rust scope unchanged; test jobs skipped" + exit 0 + fi if [ "${{ needs.test_gateway.result }}" != "success" ] || \ [ "${{ needs.test_data.result }}" != "success" ] || \ [ "${{ needs.check_data_features.result }}" != "success" ] || \ @@ -499,6 +542,8 @@ jobs: data_db_smoke_postgres: name: Data DB Smoke (Postgres) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest services: postgres: @@ -519,6 +564,8 @@ jobs: - name: Install Rust toolchain uses: dtolnay/rust-toolchain@4360b52568e2003a75bf9bc1d59f33a8e3fc893c # stable + with: + toolchain: 1.95.0 - name: Show Rust toolchain run: rustup show active-toolchain @@ -586,11 +633,16 @@ jobs: name: Data DB Smoke runs-on: ubuntu-latest needs: + - changes - data_db_smoke_postgres if: ${{ always() }} steps: - name: Verify database smoke jobs run: | + if [ "${{ needs.changes.outputs.rust }}" != "true" ]; then + echo "Rust scope unchanged; database smoke jobs skipped" + exit 0 + fi if [ "${{ needs.data_db_smoke_postgres.result }}" != "success" ]; then echo "Data DB smoke failed" exit 1 @@ -600,6 +652,7 @@ jobs: name: check runs-on: ubuntu-latest needs: + - changes - fmt - clippy - test @@ -609,11 +662,25 @@ jobs: steps: - name: Verify required jobs run: | - if [ "${{ needs.fmt.result }}" != "success" ] || \ - [ "${{ needs.clippy.result }}" != "success" ] || \ - [ "${{ needs.test.result }}" != "success" ] || \ - [ "${{ needs.data_db_smoke.result }}" != "success" ] || \ - [ "${{ needs.shell_security.result }}" != "success" ]; then + rust="${{ needs.changes.outputs.rust }}" + shell="${{ needs.changes.outputs.shell }}" + + if [ "$rust" != "true" ] && [ "$shell" != "true" ]; then + echo "No Rust or shell scope changed" + exit 0 + fi + + if [ "$rust" = "true" ] && { + [ "${{ needs.fmt.result }}" != "success" ] || + [ "${{ needs.clippy.result }}" != "success" ] || + [ "${{ needs.test.result }}" != "success" ] || + [ "${{ needs.data_db_smoke.result }}" != "success" ]; + }; then + echo "Rust CI failed" + exit 1 + fi + + if [ "$shell" = "true" ] && [ "${{ needs.shell_security.result }}" != "success" ]; then echo "Rust CI failed" exit 1 fi diff --git a/Cargo.lock b/Cargo.lock index 991a6330a..ab01ea5dc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -429,11 +429,14 @@ dependencies = [ "aether-runtime", "aether-runtime-state", "aether-testkit", + "aether-tunnel", + "arc-swap", "async-stream", "axum", "futures-util", "http", "reqwest 0.12.28", + "rustls", "serde", "serde_json", "sha2", @@ -672,7 +675,6 @@ name = "aether-tunnel" version = "0.3.17" dependencies = [ "aether-contracts", - "aether-gateway", "aether-gateway-tunnel", "aether-http", "aether-runtime", diff --git a/Cargo.toml b/Cargo.toml index 96272433b..9bc1273aa 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -92,6 +92,7 @@ aether-usage-core = { path = "crates/aether-usage/core" } aether-usage-runtime = { path = "crates/aether-usage/runtime" } aether-video-tasks-core = { path = "crates/aether-video-tasks-core" } aether-gateway = { path = "apps/aether-gateway" } +aether-tunnel = { path = "apps/aether-tunnel" } aether-http = { path = "crates/aether-http" } aether-runtime = { path = "crates/aether-runtime/base" } aether-testkit = { path = "crates/aether-testing/testkit" } diff --git a/apps/aether-gateway/src/backup/executor.rs b/apps/aether-gateway/src/backup/executor.rs index a7bc20baf..5b518b916 100644 --- a/apps/aether-gateway/src/backup/executor.rs +++ b/apps/aether-gateway/src/backup/executor.rs @@ -799,6 +799,7 @@ mod tests { use aes_gcm::aead::{Aead, AeadCore, KeyInit, OsRng, Payload}; use aes_gcm::Aes256Gcm; use aether_crypto::DEVELOPMENT_ENCRYPTION_KEY; + use base64::Engine as _; use bytes::Bytes; use chrono::{DateTime, Utc}; use serde_json::json; @@ -1243,8 +1244,18 @@ mod tests { assert_eq!(restored.key_id, None); assert_eq!(restored.export_version.as_deref(), Some("2.3")); + // 17 个互不相同的合法 base64-32 字节直接密钥:本段只验证“legacy 候选 >16 → TooManyLegacyKeys”, + // 不测口令强度、不解密。直接密钥走 decode_direct_fernet_key(生产已支持路径),跳过 PBKDF2, + // 避免本用例为计数语义再付 17×10 万次迭代;上半段 DEVELOPMENT_ENCRYPTION_KEY 真实 v1 兼容 + // 与 wrong-legacy-secret 派生路径保持不变。 let too_many: Vec<_> = (0..17) - .map(|index| BackupDecryptionKey::historical(format!("legacy-{index}")).unwrap()) + .map(|index| { + let mut material = [0u8; 32]; + material[0] = index as u8 + 1; + material[31] = index as u8 + 1; + let secret = base64::engine::general_purpose::STANDARD.encode(material); + BackupDecryptionKey::historical(secret).unwrap() + }) .collect(); assert!(matches!( restore_backup_json( diff --git a/apps/aether-gateway/src/dispatch/pool_scheduler.rs b/apps/aether-gateway/src/dispatch/pool_scheduler.rs index 5a86b5590..dcb60a8ee 100644 --- a/apps/aether-gateway/src/dispatch/pool_scheduler.rs +++ b/apps/aether-gateway/src/dispatch/pool_scheduler.rs @@ -5165,15 +5165,6 @@ mod tests { )) } - fn provider_catalog_credential_state() -> AppState { - AppState::new() - .expect("credential state should build") - .with_data_state_for_tests( - GatewayDataState::disabled() - .with_encryption_key_for_tests(aether_crypto::DEVELOPMENT_ENCRYPTION_KEY), - ) - } - fn large_pool_fixture( key_count: usize, provider_config: Option, @@ -5224,18 +5215,12 @@ mod tests { ) .expect("endpoint transport should build"); - let credential_state = provider_catalog_credential_state(); + // 这些用例只验证池扫描、跳过计数和游标预算,不会发起请求或读取凭据。 + // 留空凭据可跳过无关的 Fernet 加解密,同时避免复用绑定密文破坏 key_id AAD。 let mut keys = Vec::with_capacity(key_count); let mut rows = Vec::with_capacity(key_count); for index in 0..key_count { let key_id = format!("key-{index:05}"); - let encrypted_api_key = credential_state - .seal_provider_catalog_key_api_key( - "provider-pool", - &key_id, - &format!("secret-{index}"), - ) - .expect("api key should encrypt"); let mut key = StoredProviderCatalogKey::new( key_id.clone(), "provider-pool".to_string(), @@ -5247,7 +5232,7 @@ mod tests { .expect("key should build") .with_transport_fields( Some(json!(["openai:chat"])), - encrypted_api_key, + None, None, None, None, @@ -5382,10 +5367,8 @@ mod tests { .expect("endpoint transport should build") } + /// 这些测试只检查池调度状态,不涉及凭据解密,因此不构造无关的密文。 fn sample_codex_pool_key(provider_id: &str, key_id: &str) -> StoredProviderCatalogKey { - let encrypted_api_key = provider_catalog_credential_state() - .seal_provider_catalog_key_api_key(provider_id, key_id, &format!("secret-{key_id}")) - .expect("api key should encrypt"); let mut key = StoredProviderCatalogKey::new( key_id.to_string(), provider_id.to_string(), @@ -5397,7 +5380,7 @@ mod tests { .expect("key should build") .with_transport_fields( Some(json!(["openai:responses"])), - encrypted_api_key, + None, None, None, Some(json!({"openai:responses": 1})), diff --git a/apps/aether-gateway/src/tests/control/admin/endpoints/keys.rs b/apps/aether-gateway/src/tests/control/admin/endpoints/keys.rs index 1043eb702..7b9afb6dd 100644 --- a/apps/aether-gateway/src/tests/control/admin/endpoints/keys.rs +++ b/apps/aether-gateway/src/tests/control/admin/endpoints/keys.rs @@ -44,21 +44,11 @@ where F: FnOnce() -> Fut + Send + 'static, Fut: std::future::Future + 'static, { - let handle = std::thread::Builder::new() - .name(test_name.to_string()) - .stack_size(PROVIDER_KEYS_TEST_STACK_BYTES) - .spawn(move || { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("test runtime should build"); - runtime.block_on(make_future()); - }) - .expect("provider keys test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack( + test_name, + PROVIDER_KEYS_TEST_STACK_BYTES, + make_future, + ); } struct SummaryNullingProviderCatalogReadRepository { diff --git a/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs b/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs index 37beb9b9a..2f9b61bd5 100644 --- a/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs +++ b/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs @@ -59,21 +59,11 @@ where F: FnOnce() -> Fut + Send + 'static, Fut: std::future::Future + 'static, { - let handle = std::thread::Builder::new() - .name(test_name.to_string()) - .stack_size(PROVIDER_QUOTA_TEST_STACK_BYTES) - .spawn(move || { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("test runtime should build"); - runtime.block_on(make_future()); - }) - .expect("provider quota test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack( + test_name, + PROVIDER_QUOTA_TEST_STACK_BYTES, + make_future, + ); } #[tokio::test] diff --git a/apps/aether-gateway/src/tests/control/admin/oauth.rs b/apps/aether-gateway/src/tests/control/admin/oauth.rs index 50e9cfc95..fea3ca978 100644 --- a/apps/aether-gateway/src/tests/control/admin/oauth.rs +++ b/apps/aether-gateway/src/tests/control/admin/oauth.rs @@ -56,21 +56,11 @@ where F: FnOnce() -> Fut + Send + 'static, Fut: std::future::Future + 'static, { - let handle = std::thread::Builder::new() - .name(test_name.to_string()) - .stack_size(ADMIN_OAUTH_TEST_STACK_BYTES) - .spawn(move || { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("test runtime should build"); - runtime.block_on(make_future()); - }) - .expect("admin oauth test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack( + test_name, + ADMIN_OAUTH_TEST_STACK_BYTES, + make_future, + ); } fn decrypt_persisted_provider_api_key(key: &StoredProviderCatalogKey) -> String { diff --git a/apps/aether-gateway/src/tests/control/admin/provider_ops.rs b/apps/aether-gateway/src/tests/control/admin/provider_ops.rs index ccc2fee1f..a9b2f5647 100644 --- a/apps/aether-gateway/src/tests/control/admin/provider_ops.rs +++ b/apps/aether-gateway/src/tests/control/admin/provider_ops.rs @@ -68,21 +68,11 @@ where F: FnOnce() -> Fut + Send + 'static, Fut: std::future::Future + 'static, { - let handle = std::thread::Builder::new() - .name(test_name.to_string()) - .stack_size(PROVIDER_OPS_TEST_STACK_BYTES) - .spawn(move || { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("test runtime should build"); - runtime.block_on(make_future()); - }) - .expect("provider ops test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack( + test_name, + PROVIDER_OPS_TEST_STACK_BYTES, + make_future, + ); } async fn start_managed_redis_or_skip() -> Option { diff --git a/apps/aether-gateway/src/tests/control/admin/provider_query.rs b/apps/aether-gateway/src/tests/control/admin/provider_query.rs index 801385e34..336063114 100644 --- a/apps/aether-gateway/src/tests/control/admin/provider_query.rs +++ b/apps/aether-gateway/src/tests/control/admin/provider_query.rs @@ -36,21 +36,11 @@ where F: FnOnce() -> Fut + Send + 'static, Fut: std::future::Future + 'static, { - let handle = std::thread::Builder::new() - .name(test_name.to_string()) - .stack_size(PROVIDER_QUERY_TEST_STACK_BYTES) - .spawn(move || { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("test runtime should build"); - runtime.block_on(make_future()); - }) - .expect("provider query test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack( + test_name, + PROVIDER_QUERY_TEST_STACK_BYTES, + make_future, + ); } fn crc32(data: &[u8]) -> u32 { diff --git a/apps/aether-gateway/src/tests/control/admin/security.rs b/apps/aether-gateway/src/tests/control/admin/security.rs index 2df981435..71b469cf3 100644 --- a/apps/aether-gateway/src/tests/control/admin/security.rs +++ b/apps/aether-gateway/src/tests/control/admin/security.rs @@ -239,11 +239,25 @@ async fn send_admin_security_request( method: reqwest::Method, path: &str, body: Option, +) -> (StatusCode, serde_json::Value, usize) { + let path = path.to_string(); + crate::tests::run_async_test_on_large_stack_with_result( + "admin-security-router-request", + 16 * 1024 * 1024, + move || send_admin_security_request_on_large_stack(gateway, method, path, body), + ) +} + +async fn send_admin_security_request_on_large_stack( + gateway: Router, + method: reqwest::Method, + path: String, + body: Option, ) -> (StatusCode, serde_json::Value, usize) { let upstream_hits = Arc::new(Mutex::new(0usize)); let upstream_hits_clone = Arc::clone(&upstream_hits); let upstream = Router::new().route( - path, + &path, any(move |_request: Request| { let upstream_hits_inner = Arc::clone(&upstream_hits_clone); async move { @@ -253,26 +267,52 @@ async fn send_admin_security_request( }), ); - let (upstream_url, upstream_handle) = start_server(upstream).await; - let (gateway_url, gateway_handle) = start_server(gateway).await; + let (_upstream_url, upstream_handle) = start_server(upstream).await; - let client = reqwest::Client::new(); - let mut request = client - .request(method, format!("{gateway_url}{path}")) + // 这些用例只验证本地安全路由和“不得转发”断言,不需要为 Gateway + // 再启动一个 TCP listener;send_request 会补齐 ConnectInfo,仍经过完整 Router。 + let mut request_builder = Request::builder() + .method(method.as_str()) + .uri(&path) .header(crate::constants::GATEWAY_HEADER, "rust-phase3b") .header(TRUSTED_ADMIN_USER_ID_HEADER, "admin-user-123") .header(TRUSTED_ADMIN_USER_ROLE_HEADER, "admin") .header(TRUSTED_ADMIN_SESSION_ID_HEADER, "session-123"); if let Some(body) = body { - request = request.json(&body); + request_builder = request_builder.header(http::header::CONTENT_TYPE, "application/json"); + let request = request_builder + .body(Body::from(body.to_string())) + .expect("request should build"); + let response = send_request(gateway, request).await; + let status = response.status(); + let payload = response + .into_body() + .collect() + .await + .expect("response body should collect") + .to_bytes(); + let payload: serde_json::Value = + serde_json::from_slice(&payload).expect("json body should parse"); + let upstream_count = *upstream_hits.lock().expect("mutex should lock"); + upstream_handle.abort(); + return (status, payload, upstream_count); } - let response = request.send().await.expect("request should succeed"); + let request = request_builder + .body(Body::empty()) + .expect("request should build"); + let response = send_request(gateway, request).await; let status = response.status(); - let payload: serde_json::Value = response.json().await.expect("json body should parse"); + let payload = response + .into_body() + .collect() + .await + .expect("response body should collect") + .to_bytes(); + let payload: serde_json::Value = + serde_json::from_slice(&payload).expect("json body should parse"); let upstream_count = *upstream_hits.lock().expect("mutex should lock"); - gateway_handle.abort(); upstream_handle.abort(); (status, payload, upstream_count) diff --git a/apps/aether-gateway/src/tests/control/admin/system_import.rs b/apps/aether-gateway/src/tests/control/admin/system_import.rs index 022a7ece6..8ad2871b1 100644 --- a/apps/aether-gateway/src/tests/control/admin/system_import.rs +++ b/apps/aether-gateway/src/tests/control/admin/system_import.rs @@ -341,21 +341,11 @@ where F: FnOnce() -> Fut + Send + 'static, Fut: std::future::Future + 'static, { - let handle = std::thread::Builder::new() - .name(test_name.to_string()) - .stack_size(ADMIN_SYSTEM_IMPORT_TEST_STACK_BYTES) - .spawn(move || { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("test runtime should build"); - runtime.block_on(make_future()); - }) - .expect("admin system import test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack( + test_name, + ADMIN_SYSTEM_IMPORT_TEST_STACK_BYTES, + make_future, + ); } #[test] diff --git a/apps/aether-gateway/src/tests/files/mod.rs b/apps/aether-gateway/src/tests/files/mod.rs index a9c9b1e7c..0acd7fb10 100644 --- a/apps/aether-gateway/src/tests/files/mod.rs +++ b/apps/aether-gateway/src/tests/files/mod.rs @@ -37,21 +37,7 @@ where F: FnOnce() -> Fut + Send + 'static, Fut: std::future::Future + 'static, { - let handle = std::thread::Builder::new() - .name(test_name.to_string()) - .stack_size(FILES_TEST_STACK_BYTES) - .spawn(move || { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("test runtime should build"); - runtime.block_on(make_future()); - }) - .expect("files test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack(test_name, FILES_TEST_STACK_BYTES, make_future); } fn hash_api_key(value: &str) -> String { diff --git a/apps/aether-gateway/src/tests/frontdoor.rs b/apps/aether-gateway/src/tests/frontdoor.rs index 15ae3a645..32a20d3e4 100644 --- a/apps/aether-gateway/src/tests/frontdoor.rs +++ b/apps/aether-gateway/src/tests/frontdoor.rs @@ -39,21 +39,7 @@ fn run_frontdoor_async_test(name: &'static str, future: F) where F: std::future::Future + Send + 'static, { - let handle = std::thread::Builder::new() - .name(name.to_string()) - .stack_size(16 * 1024 * 1024) - .spawn(move || { - tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("frontdoor test runtime should build") - .block_on(future); - }) - .expect("large-stack frontdoor test thread should spawn"); - - if let Err(payload) = handle.join() { - std::panic::resume_unwind(payload); - } + crate::tests::run_async_test_on_large_stack(name, 16 * 1024 * 1024, || future); } fn hash_api_key(value: &str) -> String { diff --git a/apps/aether-gateway/src/tests/mod.rs b/apps/aether-gateway/src/tests/mod.rs index 15d90bcf9..4e93f587a 100644 --- a/apps/aether-gateway/src/tests/mod.rs +++ b/apps/aether-gateway/src/tests/mod.rs @@ -10,7 +10,6 @@ pub(super) use http::StatusCode; pub(super) use serde_json::json; mod ai_execute; -mod architecture; mod async_task; mod audit; mod concurrency; @@ -46,6 +45,50 @@ pub(super) async fn start_server(app: Router) -> (String, tokio::task::JoinHandl (format!("http://{addr}"), handle) } +/// 在独立的大栈线程中运行需要深调用栈的异步测试。 +/// +/// 这些测试仍保留 16 MiB 栈空间;这里只统一线程和 runtime 的启动逻辑, +/// 避免每个测试分区各自复制一份 helper,降低维护时误改测试执行语义的风险。 +pub(crate) fn run_async_test_on_large_stack( + test_name: &'static str, + stack_size: usize, + make_future: F, +) where + F: FnOnce() -> Fut + Send + 'static, + Fut: std::future::Future + 'static, +{ + run_async_test_on_large_stack_with_result(test_name, stack_size, make_future); +} + +/// 与上面的 helper 相同,但允许深栈测试返回结果,供公共请求 helper 使用。 +pub(crate) fn run_async_test_on_large_stack_with_result( + test_name: &'static str, + stack_size: usize, + make_future: F, +) -> R +where + F: FnOnce() -> Fut + Send + 'static, + Fut: std::future::Future + 'static, + R: Send + 'static, +{ + let handle = std::thread::Builder::new() + .name(test_name.to_string()) + .stack_size(stack_size) + .spawn(move || { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("test runtime should build"); + runtime.block_on(make_future()) + }) + .expect("large-stack test thread should spawn"); + + match handle.join() { + Ok(result) => result, + Err(payload) => std::panic::resume_unwind(payload), + } +} + pub(super) const OPERATIONAL_ADMIN_DEVICE_ID: &str = "device-operational-admin"; pub(super) async fn start_authenticated_operational_server( diff --git a/apps/aether-gateway/src/tests/architecture/admin_billing.rs b/apps/aether-gateway/tests/architecture/admin_billing.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/admin_billing.rs rename to apps/aether-gateway/tests/architecture/admin_billing.rs diff --git a/apps/aether-gateway/src/tests/architecture/admin_model.rs b/apps/aether-gateway/tests/architecture/admin_model.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/admin_model.rs rename to apps/aether-gateway/tests/architecture/admin_model.rs diff --git a/apps/aether-gateway/src/tests/architecture/admin_observability.rs b/apps/aether-gateway/tests/architecture/admin_observability.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/admin_observability.rs rename to apps/aether-gateway/tests/architecture/admin_observability.rs diff --git a/apps/aether-gateway/src/tests/architecture/admin_provider.rs b/apps/aether-gateway/tests/architecture/admin_provider.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/admin_provider.rs rename to apps/aether-gateway/tests/architecture/admin_provider.rs diff --git a/apps/aether-gateway/src/tests/architecture/admin_shared.rs b/apps/aether-gateway/tests/architecture/admin_shared.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/admin_shared.rs rename to apps/aether-gateway/tests/architecture/admin_shared.rs diff --git a/apps/aether-gateway/src/tests/architecture/admin_system.rs b/apps/aether-gateway/tests/architecture/admin_system.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/admin_system.rs rename to apps/aether-gateway/tests/architecture/admin_system.rs diff --git a/apps/aether-gateway/src/tests/architecture/admin_users.rs b/apps/aether-gateway/tests/architecture/admin_users.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/admin_users.rs rename to apps/aether-gateway/tests/architecture/admin_users.rs diff --git a/apps/aether-gateway/src/tests/architecture/ai_serving.rs b/apps/aether-gateway/tests/architecture/ai_serving.rs similarity index 99% rename from apps/aether-gateway/src/tests/architecture/ai_serving.rs rename to apps/aether-gateway/tests/architecture/ai_serving.rs index fa71ab5c7..c429de37c 100644 --- a/apps/aether-gateway/src/tests/architecture/ai_serving.rs +++ b/apps/aether-gateway/tests/architecture/ai_serving.rs @@ -1194,7 +1194,7 @@ fn ai_serving_planner_separates_local_candidate_resolution_from_ranking() { let candidate_resolution = read_workspace_file("apps/aether-gateway/src/ai_serving/planner/candidate_resolution.rs"); - let ranking_call = candidate_resolution + candidate_resolution .find("rank_eligible_local_execution_candidates(") .expect("candidate_resolution.rs should call core-backed local candidate ranking"); assert!( @@ -5045,7 +5045,9 @@ fn retired_api_format_occurrences_are_whitelisted() { .expect("file should be under workspace root") .to_string_lossy() .replace('\\', "/"); - if relative == "apps/aether-gateway/src/tests/architecture/ai_serving.rs" { + if relative == "apps/aether-gateway/tests/architecture/ai_serving.rs" + || relative == "apps/aether-gateway/src/tests/architecture/ai_serving.rs" + { continue; } diff --git a/apps/aether-gateway/src/tests/architecture/mod.rs b/apps/aether-gateway/tests/architecture/mod.rs similarity index 89% rename from apps/aether-gateway/src/tests/architecture/mod.rs rename to apps/aether-gateway/tests/architecture/mod.rs index 374e1074b..57ddab353 100644 --- a/apps/aether-gateway/src/tests/architecture/mod.rs +++ b/apps/aether-gateway/tests/architecture/mod.rs @@ -1,7 +1,9 @@ use std::fs; use std::path::{Path, PathBuf}; -pub(super) fn collect_rust_files(root: &Path, files: &mut Vec) { +// 架构守卫在独立 integration test 中是顶层模块;helper 统一 pub(crate), +// 子模块经 `use super::*` / `use super::{...}` 访问(与原 lib 内布局一致)。 +pub(crate) fn collect_rust_files(root: &Path, files: &mut Vec) { for entry in fs::read_dir(root).expect("directory should be readable") { let entry = entry.expect("directory entry should be readable"); let path = entry.path(); @@ -15,7 +17,7 @@ pub(super) fn collect_rust_files(root: &Path, files: &mut Vec) { } } -pub(super) fn assert_no_sqlx_queries(root_relative_path: &str) { +pub(crate) fn assert_no_sqlx_queries(root_relative_path: &str) { let root = Path::new(env!("CARGO_MANIFEST_DIR")).join(root_relative_path); let mut files = Vec::new(); collect_rust_files(&root, &mut files); @@ -80,7 +82,7 @@ fn sql_pool_scan_distinguishes_pool_types_from_repository_names() { )); } -pub(super) fn assert_no_sensitive_log_patterns(root_relative_path: &str, patterns: &[&str]) { +pub(crate) fn assert_no_sensitive_log_patterns(root_relative_path: &str, patterns: &[&str]) { let root = Path::new(env!("CARGO_MANIFEST_DIR")).join(root_relative_path); let mut files = Vec::new(); collect_rust_files(&root, &mut files); @@ -109,7 +111,7 @@ pub(super) fn assert_no_sensitive_log_patterns(root_relative_path: &str, pattern ); } -pub(super) fn assert_no_module_dependency_patterns(root_relative_path: &str, patterns: &[&str]) { +pub(crate) fn assert_no_module_dependency_patterns(root_relative_path: &str, patterns: &[&str]) { let root = Path::new(env!("CARGO_MANIFEST_DIR")).join(root_relative_path); let mut files = Vec::new(); collect_rust_files(&root, &mut files); @@ -138,14 +140,14 @@ pub(super) fn assert_no_module_dependency_patterns(root_relative_path: &str, pat ); } -pub(super) fn workspace_file_exists(root_relative_path: &str) -> bool { +pub(crate) fn workspace_file_exists(root_relative_path: &str) -> bool { Path::new(env!("CARGO_MANIFEST_DIR")) .join("../..") .join(root_relative_path) .exists() } -pub(super) fn workspace_files_with_extension( +pub(crate) fn workspace_files_with_extension( root_relative_path: &str, extension: &str, ) -> Vec { @@ -162,7 +164,7 @@ pub(super) fn workspace_files_with_extension( files } -pub(super) fn collect_workspace_rust_files(root_relative_path: &str) -> Vec { +pub(crate) fn collect_workspace_rust_files(root_relative_path: &str) -> Vec { let root = Path::new(env!("CARGO_MANIFEST_DIR")) .join("../..") .join(root_relative_path); @@ -172,7 +174,7 @@ pub(super) fn collect_workspace_rust_files(root_relative_path: &str) -> Vec String { +pub(crate) fn read_workspace_file(path: &str) -> String { let workspace_root = Path::new(env!("CARGO_MANIFEST_DIR")) .join("../..") .canonicalize() @@ -180,7 +182,7 @@ pub(super) fn read_workspace_file(path: &str) -> String { fs::read_to_string(workspace_root.join(path)).expect("source file should be readable") } -pub(super) fn read_workspace_module_tree(path: &str) -> String { +pub(crate) fn read_workspace_module_tree(path: &str) -> String { let workspace_root = Path::new(env!("CARGO_MANIFEST_DIR")) .join("../..") .canonicalize() diff --git a/apps/aether-gateway/src/tests/architecture/runtime_and_security.rs b/apps/aether-gateway/tests/architecture/runtime_and_security.rs similarity index 99% rename from apps/aether-gateway/src/tests/architecture/runtime_and_security.rs rename to apps/aether-gateway/tests/architecture/runtime_and_security.rs index 8772263c8..4cff5b71c 100644 --- a/apps/aether-gateway/src/tests/architecture/runtime_and_security.rs +++ b/apps/aether-gateway/tests/architecture/runtime_and_security.rs @@ -1,4 +1,4 @@ -use std::path::{Path, PathBuf}; +use std::path::Path; use super::*; diff --git a/apps/aether-gateway/src/tests/architecture/sql_and_data.rs b/apps/aether-gateway/tests/architecture/sql_and_data.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/sql_and_data.rs rename to apps/aether-gateway/tests/architecture/sql_and_data.rs diff --git a/apps/aether-gateway/src/tests/architecture/usage.rs b/apps/aether-gateway/tests/architecture/usage.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/usage.rs rename to apps/aether-gateway/tests/architecture/usage.rs diff --git a/apps/aether-gateway/src/tests/architecture/workspace_tiers.rs b/apps/aether-gateway/tests/architecture/workspace_tiers.rs similarity index 100% rename from apps/aether-gateway/src/tests/architecture/workspace_tiers.rs rename to apps/aether-gateway/tests/architecture/workspace_tiers.rs diff --git a/apps/aether-gateway/tests/architecture_guard.rs b/apps/aether-gateway/tests/architecture_guard.rs new file mode 100644 index 000000000..ebe87041c --- /dev/null +++ b/apps/aether-gateway/tests/architecture_guard.rs @@ -0,0 +1,5 @@ +//! 架构守卫独立测试目标。 +//! +//! 从 lib 的 `cfg(test)` 巨型编译单元迁出:只做源码/manifest 字符串断言, +//! 不启动 AppState、不依赖 gateway 私有类型,用于压低 lib test 编译面与 rustc 峰值。 +mod architecture; diff --git a/apps/aether-tunnel/Cargo.toml b/apps/aether-tunnel/Cargo.toml index 330b7127a..480f4f267 100644 --- a/apps/aether-tunnel/Cargo.toml +++ b/apps/aether-tunnel/Cargo.toml @@ -46,5 +46,4 @@ webpki-roots = "0.26" uuid.workspace = true [dev-dependencies] -aether-gateway = { workspace = true, features = ["testkit"] } tokio = { version = "1", features = ["test-util"] } diff --git a/apps/aether-tunnel/src/config.rs b/apps/aether-tunnel/src/config.rs index 5c45c495b..0cf824188 100644 --- a/apps/aether-tunnel/src/config.rs +++ b/apps/aether-tunnel/src/config.rs @@ -885,7 +885,7 @@ impl Config { Ok(Duration::from_millis(self.tunnel_connect_timeout_ms)) } - pub fn tunnel_ip_family(&self) -> crate::egress_proxy::IpFamily { + pub(crate) fn tunnel_ip_family(&self) -> crate::egress_proxy::IpFamily { if self.tunnel_ipv4_only { crate::egress_proxy::IpFamily::Ipv4Only } else if self.tunnel_ipv6_only { diff --git a/apps/aether-tunnel/src/hardware.rs b/apps/aether-tunnel/src/hardware.rs index d0e9a25fc..01192c96c 100644 --- a/apps/aether-tunnel/src/hardware.rs +++ b/apps/aether-tunnel/src/hardware.rs @@ -116,6 +116,13 @@ impl RuntimeResourceMonitor { } } +impl Default for RuntimeResourceMonitor { + fn default() -> Self { + // 默认构造与显式 new 保持一致,便于库目标和二进制目标共用监控器。 + Self::new() + } +} + /// Collect hardware information and estimate max concurrency. /// /// Should be called once at startup -- hardware does not change at runtime. diff --git a/apps/aether-tunnel/src/lib.rs b/apps/aether-tunnel/src/lib.rs new file mode 100644 index 000000000..98ef478a2 --- /dev/null +++ b/apps/aether-tunnel/src/lib.rs @@ -0,0 +1,17 @@ +#![allow(clippy::large_enum_variant)] + +// Tunnel 的运行模块作为库暴露给独立集成测试使用;生产二进制仍由 +// src/main.rs 负责命令行解析,避免端到端测试把 Gateway dev-dependency +// 带进 Workspace Rest 的默认测试目标。 +pub mod app; +pub mod config; +pub mod egress_proxy; +pub mod hardware; +mod net; +pub mod registration; +pub mod runtime; +pub mod setup; +pub mod state; +pub mod target_filter; +pub mod tunnel; +pub mod upstream_client; diff --git a/apps/aether-tunnel/src/main.rs b/apps/aether-tunnel/src/main.rs index d2b01e23d..9ea67bf51 100644 --- a/apps/aether-tunnel/src/main.rs +++ b/apps/aether-tunnel/src/main.rs @@ -1,20 +1,8 @@ #![allow(clippy::large_enum_variant)] -mod app; -mod config; -mod egress_proxy; -mod hardware; -mod net; -mod registration; -mod runtime; -mod setup; -mod state; -mod target_filter; -mod tunnel; -mod upstream_client; - use std::path::PathBuf; +use aether_tunnel::{app, config, setup}; use clap::{parser::ValueSource, CommandFactory, FromArgMatches, Parser}; use config::{Config, ServerEntry, TunnelSecurity}; diff --git a/apps/aether-tunnel/src/setup/mod.rs b/apps/aether-tunnel/src/setup/mod.rs index 3c9628c7c..0437863fb 100644 --- a/apps/aether-tunnel/src/setup/mod.rs +++ b/apps/aether-tunnel/src/setup/mod.rs @@ -1,5 +1,5 @@ -pub(crate) mod service; +pub mod service; mod tui; -pub(crate) mod upgrade; +pub mod upgrade; pub use self::tui::{run, SetupOutcome}; diff --git a/apps/aether-tunnel/src/state.rs b/apps/aether-tunnel/src/state.rs index 26b54a42d..df5959e28 100644 --- a/apps/aether-tunnel/src/state.rs +++ b/apps/aether-tunnel/src/state.rs @@ -226,6 +226,13 @@ impl TunnelRequestMetrics { } } +impl Default for TunnelRequestMetrics { + fn default() -> Self { + // 指标初始值全部为零,Default 与现有 new 语义完全一致。 + Self::new() + } +} + const RECENT_TUNNEL_ERROR_CAPACITY: usize = 64; const TUNNEL_ERROR_CATEGORY_MAX_CHARS: usize = 48; const TUNNEL_ERROR_MESSAGE_MAX_CHARS: usize = 320; @@ -534,6 +541,13 @@ impl TunnelMetrics { } } +impl Default for TunnelMetrics { + fn default() -> Self { + // 保留 recent_errors 的容量初始化,避免 Default 改变错误环形缓存行为。 + Self::new() + } +} + fn now_unix_secs() -> u64 { now_unix_ms() / 1_000 } diff --git a/apps/aether-tunnel/src/tunnel/mod.rs b/apps/aether-tunnel/src/tunnel/mod.rs index 810899bad..b0c46986f 100644 --- a/apps/aether-tunnel/src/tunnel/mod.rs +++ b/apps/aether-tunnel/src/tunnel/mod.rs @@ -1,10 +1,10 @@ pub mod client; -pub mod dispatcher; -pub mod heartbeat; +mod dispatcher; +mod heartbeat; pub mod protocol; -pub mod stream_handler; +mod stream_handler; mod task; -pub mod writer; +mod writer; use std::sync::Arc; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; @@ -230,34 +230,10 @@ fn mix_u64(mut x: u64) -> u64 { #[cfg(test)] mod tests { - use std::sync::atomic::AtomicU64; - use std::sync::{Arc, Once}; - use std::time::{Duration, SystemTime, UNIX_EPOCH}; - - use aether_contracts::tunnel::{ - sign_tunnel_relay_request, tunnel_relay_payload_digest, TUNNEL_RELAY_AUTH_NONCE_HEADER, - TUNNEL_RELAY_AUTH_PAYLOAD_HEADER, TUNNEL_RELAY_AUTH_SENDER_HEADER, - TUNNEL_RELAY_AUTH_SIGNATURE_HEADER, TUNNEL_RELAY_AUTH_TIMESTAMP_HEADER, - TUNNEL_RELAY_OWNER_INSTANCE_HEADER, - }; - use aether_gateway::{build_router_with_state, AppState as GatewayAppState}; - use arc_swap::ArcSwap; - use axum::Router; - use reqwest::StatusCode; - use tokio::sync::watch; - - use crate::config::Config; - use crate::registration::client::AetherClient; - use crate::runtime::DynamicConfig; - use crate::state::{ - AppState as TunnelAppState, ServerContext, TunnelMetrics, TunnelRequestMetrics, - }; - use crate::target_filter::DnsCache; - use crate::tunnel::protocol; - use crate::upstream_client; + use std::time::Duration; use super::{ - compute_reconnect_cap_ms, compute_reconnect_delay, compute_startup_stagger, run, + compute_reconnect_cap_ms, compute_reconnect_delay, compute_startup_stagger, MAX_STARTUP_STAGGER_MS, RECONNECT_PROBE_MAX_DELAY_MS, STARTUP_STAGGER_STEP_MS, }; @@ -299,480 +275,4 @@ mod tests { let d = compute_reconnect_delay(500, 45_000, 100, 12345); assert!(d <= Duration::from_millis(RECONNECT_PROBE_MAX_DELAY_MS)); } - - #[tokio::test] - async fn tunnel_reconnects_after_gateway_restart() { - ensure_rustls_provider(); - - let gateway_port = reserve_local_port().expect("gateway port should reserve"); - let gateway_base_url = format!("http://127.0.0.1:{gateway_port}"); - let (gateway_state, mut gateway_handle) = start_gateway_on_port(gateway_port) - .await - .expect("gateway should start"); - - let mut tunnel_config = sample_config(&gateway_base_url); - tunnel_config.tunnel_security = crate::config::TunnelSecurity::NonTlsRequired; - tunnel_config.tunnel_encryption_key = - Some("BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=".to_string()); - let state = sample_state(tunnel_config); - let server = sample_server(&state, "node-recovery"); - let (shutdown_tx, shutdown_rx) = watch::channel(false); - let tunnel_task = tokio::spawn({ - let state = Arc::clone(&state); - let server = Arc::clone(&server); - let (_drain_tx, drain_rx) = watch::channel(false); - async move { - run(&state, &server, 0, shutdown_rx, drain_rx).await; - } - }); - - wait_until_relay_status( - &gateway_base_url, - "node-recovery", - StatusCode::GATEWAY_TIMEOUT, - ) - .await; - - gateway_handle.abort(); - let _ = (&mut gateway_handle).await; - assert_eq!(gateway_state.force_close_all_tunnel_proxies(), 1); - - let (_restarted_gateway_state, restarted_gateway_handle) = - start_gateway_on_port_retry(gateway_port) - .await - .expect("gateway should restart on fixed port"); - gateway_handle = restarted_gateway_handle; - - wait_until_relay_status( - &gateway_base_url, - "node-recovery", - StatusCode::GATEWAY_TIMEOUT, - ) - .await; - - assert!(server.tunnel_metrics.snapshot().connect_successes >= 2); - let _ = shutdown_tx.send(true); - tokio::time::timeout(Duration::from_secs(5), tunnel_task) - .await - .expect("tunnel task should stop") - .expect("tunnel task should join"); - gateway_handle.abort(); - } - - async fn wait_until_relay_status(gateway_base_url: &str, node_id: &str, expected: StatusCode) { - let deadline = tokio::time::Instant::now() + Duration::from_secs(10); - let mut last_observed = None::; - loop { - if let Some((status, body)) = probe_relay_status(gateway_base_url, node_id).await { - last_observed = Some(format!("{status} body={body}")); - if status == expected { - return; - } - } - assert!( - tokio::time::Instant::now() < deadline, - "relay status did not become {expected} within timeout; last={:?}", - last_observed - ); - tokio::time::sleep(Duration::from_millis(25)).await; - } - } - - async fn probe_relay_status( - gateway_base_url: &str, - node_id: &str, - ) -> Option<(StatusCode, String)> { - let response = relay_response(gateway_base_url, node_id, relay_probe_envelope()).await?; - let status = response.status(); - let body = response.text().await.unwrap_or_default(); - Some((status, body)) - } - - async fn relay_response( - gateway_base_url: &str, - node_id: &str, - payload: Vec, - ) -> Option { - let timestamp = SystemTime::now() - .duration_since(UNIX_EPOCH) - .expect("test clock should be after epoch") - .as_secs(); - let nonce = uuid::Uuid::new_v4().simple().to_string(); - let digest = tunnel_relay_payload_digest(&payload, &[]); - let signature = sign_tunnel_relay_request( - b"tunnel-reconnect-test-secret-at-least-32-bytes", - "tunnel-reconnect-test-client", - "tunnel-reconnect-test-gateway", - node_id, - "", - false, - timestamp, - &nonce, - &digest, - ); - reqwest::Client::new() - .post(format!( - "{gateway_base_url}/api/internal/tunnel/relay/{node_id}" - )) - .header("content-type", "application/octet-stream") - .header( - TUNNEL_RELAY_AUTH_SENDER_HEADER, - "tunnel-reconnect-test-client", - ) - .header( - TUNNEL_RELAY_OWNER_INSTANCE_HEADER, - "tunnel-reconnect-test-gateway", - ) - .header(TUNNEL_RELAY_AUTH_TIMESTAMP_HEADER, timestamp) - .header(TUNNEL_RELAY_AUTH_NONCE_HEADER, nonce) - .header( - TUNNEL_RELAY_AUTH_PAYLOAD_HEADER, - digest.encode_header_value(), - ) - .header(TUNNEL_RELAY_AUTH_SIGNATURE_HEADER, signature) - .body(payload) - .send() - .await - .ok() - } - - fn relay_probe_envelope() -> Vec { - let meta = protocol::RequestMeta { - provider_id: None, - endpoint_id: None, - key_id: None, - method: "GET".to_string(), - url: "http://127.0.0.1:80/blocked".to_string(), - headers: std::collections::HashMap::new(), - stream: false, - request_timeout_ms: None, - stream_first_byte_timeout_ms: None, - timeout: 5, - follow_redirects: None, - http1_only: false, - transport_profile: None, - }; - let meta_json = - serde_json::to_vec(&meta).expect("tunnel relay probe metadata should serialize"); - let mut envelope = Vec::with_capacity(4 + meta_json.len()); - envelope.extend_from_slice(&(meta_json.len() as u32).to_be_bytes()); - envelope.extend_from_slice(&meta_json); - envelope - } - - async fn start_gateway_on_port( - port: u16, - ) -> Result<(GatewayAppState, tokio::task::JoinHandle<()>), std::io::Error> { - // The embedded gateway now fails closed when relay authentication is - // not configured. Keep this integration fixture explicitly authenticated. - static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); - let state = { - let _guard = ENV_LOCK.lock().unwrap(); - let previous_secret = std::env::var_os("AETHER_TUNNEL_RELAY_AUTH_SECRET"); - let previous_instance = std::env::var_os("AETHER_GATEWAY_INSTANCE_ID"); - std::env::set_var( - "AETHER_TUNNEL_RELAY_AUTH_SECRET", - "tunnel-reconnect-test-secret-at-least-32-bytes", - ); - std::env::set_var( - "AETHER_GATEWAY_INSTANCE_ID", - "tunnel-reconnect-test-gateway", - ); - let mut state = GatewayAppState::new().expect("gateway test state should build"); - aether_gateway::configure_test_tunnel_security( - &mut state, - "node-recovery", - "test-generation-1", - "BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=", - ); - restore_test_env("AETHER_TUNNEL_RELAY_AUTH_SECRET", previous_secret); - restore_test_env("AETHER_GATEWAY_INSTANCE_ID", previous_instance); - state - }; - let router = build_router_with_state(state.clone()); - let handle = spawn_router_on_port(port, router).await?; - Ok((state, handle)) - } - - #[tokio::test] - async fn negotiated_small_window_streams_large_responses_and_cancels_idle_upstream() { - use axum::body::{Body, Bytes}; - use axum::routing::get; - use futures_util::StreamExt; - - ensure_rustls_provider(); - let upstream_port = reserve_local_port().unwrap(); - let upstream = Router::new() - .route( - "/large", - get(|| async { Body::from(vec![b'x'; 2 * 1024 * 1024]) }), - ) - .route( - "/idle", - get(|| async { - let first = futures_util::stream::once(async { - Ok::<_, std::io::Error>(Bytes::from_static(b"data: started\n\n")) - }); - ( - [("content-type", "text/event-stream")], - Body::from_stream(first.chain(futures_util::stream::pending())), - ) - }), - ); - let upstream_task = super::task::SessionTask::new( - spawn_router_on_port(upstream_port, upstream).await.unwrap(), - ); - let gateway_port = reserve_local_port().unwrap(); - let gateway_url = format!("http://127.0.0.1:{gateway_port}"); - let (_, gateway_task) = start_gateway_on_port(gateway_port).await.unwrap(); - let gateway_task = super::task::SessionTask::new(gateway_task); - let mut config = sample_config(&gateway_url); - config.tunnel_security = crate::config::TunnelSecurity::NonTlsRequired; - config.tunnel_encryption_key = Some("BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=".into()); - config.tunnel_stream_initial_window_bytes = 512 * 1024; - config.tunnel_drain_deadline_ms = 100; - config.allow_private_targets = true; - config.allowed_ports.push(upstream_port); - let state = sample_state(config); - let server = sample_server(&state, "node-recovery"); - let (shutdown_tx, shutdown_rx) = watch::channel(false); - let (_drain_tx, drain_rx) = watch::channel(false); - let tunnel_task = super::task::SessionTask::new(tokio::spawn({ - let state = Arc::clone(&state); - let server = Arc::clone(&server); - async move { - run(&state, &server, 0, shutdown_rx, drain_rx).await; - } - })); - wait_until_relay_status(&gateway_url, "node-recovery", StatusCode::GATEWAY_TIMEOUT).await; - - let envelope = |path: &str| { - let mut meta: protocol::RequestMeta = - serde_json::from_slice(&relay_probe_envelope()[4..]).unwrap(); - meta.url = format!("http://127.0.0.1:{upstream_port}/{path}"); - meta.stream = true; - meta.timeout = 10; - meta.stream_first_byte_timeout_ms = Some(10_000); - let encoded = serde_json::to_vec(&meta).unwrap(); - let mut result = (encoded.len() as u32).to_be_bytes().to_vec(); - result.extend(encoded); - result - }; - let response = relay_response(&gateway_url, "node-recovery", envelope("large")) - .await - .unwrap(); - assert_eq!(response.status(), StatusCode::OK); - let body = tokio::time::timeout(Duration::from_secs(10), response.bytes()) - .await - .unwrap() - .unwrap(); - assert_eq!(body.len(), 2 * 1024 * 1024); - assert!(body.iter().all(|byte| *byte == b'x')); - - let mut response = relay_response(&gateway_url, "node-recovery", envelope("idle")) - .await - .unwrap(); - assert_eq!( - response.chunk().await.unwrap().unwrap(), - "data: started\n\n" - ); - drop(response); - tokio::time::timeout(Duration::from_secs(3), async { - while server - .active_connections - .load(std::sync::atomic::Ordering::Acquire) - != 0 - { - tokio::task::yield_now().await; - } - }) - .await - .expect("cancelled SSE must release the upstream handler"); - - let mut response = relay_response(&gateway_url, "node-recovery", envelope("idle")) - .await - .unwrap(); - assert!(response.chunk().await.unwrap().is_some()); - shutdown_tx.send(true).unwrap(); - tokio::time::timeout(Duration::from_secs(3), tunnel_task) - .await - .unwrap() - .unwrap(); - assert_eq!( - server - .active_connections - .load(std::sync::atomic::Ordering::Acquire), - 0 - ); - drop(response); - drop(gateway_task); - drop(upstream_task); - } - - fn restore_test_env(key: &str, value: Option) { - if let Some(value) = value { - std::env::set_var(key, value); - } else { - std::env::remove_var(key); - } - } - - async fn start_gateway_on_port_retry( - port: u16, - ) -> Result<(GatewayAppState, tokio::task::JoinHandle<()>), std::io::Error> { - let mut attempts = 0usize; - loop { - match start_gateway_on_port(port).await { - Ok(server) => return Ok(server), - Err(err) => { - attempts += 1; - if attempts >= 20 { - return Err(err); - } - tokio::time::sleep(Duration::from_millis(50)).await; - } - } - } - } - - async fn spawn_router_on_port( - port: u16, - app: Router, - ) -> Result, std::io::Error> { - let listener = tokio::net::TcpListener::bind(("127.0.0.1", port)).await?; - Ok(tokio::spawn(async move { - axum::serve( - listener, - app.into_make_service_with_connect_info::(), - ) - .await - .expect("gateway test server should run"); - })) - } - - fn reserve_local_port() -> Result { - let listener = std::net::TcpListener::bind("127.0.0.1:0")?; - let port = listener.local_addr()?.port(); - drop(listener); - Ok(port) - } - - fn sample_state(config: Config) -> Arc { - let config = Arc::new(config); - let dns_cache = Arc::new(DnsCache::new(Duration::from_secs(60), 128)); - let upstream_client_pool = - upstream_client::UpstreamClientPool::new(Arc::clone(&config), Arc::clone(&dns_cache)); - Arc::new(TunnelAppState { - config, - dns_cache, - upstream_client_pool, - tunnel_tls_config: Arc::new(crate::tunnel::client::build_tls_config()), - resource_monitor: Arc::new(crate::hardware::RuntimeResourceMonitor::new()), - stream_gate: None, - distributed_stream_gate: None, - }) - } - - fn sample_server(state: &Arc, node_id: &str) -> Arc { - let config = Arc::clone(&state.config); - Arc::new(ServerContext { - server_label: "gateway-owned-tunnel".to_string(), - aether_url: config.aether_url.clone(), - management_token: config.management_token.clone(), - tunnel_security: config.tunnel_security, - tunnel_encryption_key: config.tunnel_encryption_key.clone(), - node_name: config.node_name.clone(), - node_id: Arc::new(std::sync::RwLock::new(node_id.to_string())), - tunnel_generation: "test-generation-1".to_string(), - aether_client: Arc::new(AetherClient::new( - &config, - &config.aether_url, - &config.management_token, - )), - dynamic: Arc::new(ArcSwap::from_pointee(DynamicConfig::from_config(&config))), - active_connections: Arc::new(AtomicU64::new(0)), - metrics: Arc::new(TunnelRequestMetrics::new()), - tunnel_metrics: Arc::new(TunnelMetrics::new()), - }) - } - - fn sample_config(aether_url: &str) -> Config { - Config { - aether_url: aether_url.to_string(), - management_token: "token".to_string(), - public_ip: None, - node_name: "tunnel-test".to_string(), - tunnel_security: crate::config::TunnelSecurity::Off, - tunnel_encryption_key: None, - node_region: None, - heartbeat_interval: 1, - allowed_ports: vec![80, 443], - allow_private_targets: false, - aether_request_timeout_secs: 10, - aether_connect_timeout_secs: 2, - aether_pool_max_idle_per_host: 8, - aether_pool_idle_timeout_secs: 90, - aether_tcp_keepalive_secs: 60, - aether_tcp_nodelay: true, - aether_http2: true, - aether_outbound_proxy_url: None, - aether_retry_max_attempts: 1, - aether_retry_base_delay_ms: 50, - aether_retry_max_delay_ms: 100, - diagnostics_bind: None, - max_concurrent_connections: None, - max_in_flight_streams: None, - distributed_stream_limit: None, - distributed_stream_redis_url: None, - distributed_stream_redis_key_prefix: None, - distributed_stream_lease_ttl_ms: 30_000, - distributed_stream_renew_interval_ms: 10_000, - distributed_stream_command_timeout_ms: 1_000, - dns_cache_ttl_secs: 60, - dns_cache_capacity: 128, - upstream_connect_timeout_secs: 30, - upstream_pool_max_idle_per_host: 4, - upstream_pool_idle_timeout_secs: 60, - upstream_client_pool_capacity: crate::config::DEFAULT_UPSTREAM_CLIENT_POOL_CAPACITY, - upstream_tcp_keepalive_secs: 60, - upstream_tcp_nodelay: true, - upstream_proxy_url: None, - upstream_proxy_remote_dns: false, - legacy_redirect_replay_budget_bytes_ignored: None, - emit_proxy_timing_header: true, - log_level: "info".to_string(), - log_destination: crate::config::TunnelLogDestinationArg::Stdout, - log_dir: None, - log_rotation: crate::config::TunnelLogRotationArg::Daily, - log_retention_days: 7, - log_max_files: 30, - tunnel_reconnect_base_ms: 50, - tunnel_reconnect_max_ms: 250, - tunnel_ping_interval_ms: 1_000, - tunnel_max_streams: Some(8), - tunnel_profile: crate::config::TunnelProfileArg::Lite, - tunnel_stream_initial_window_bytes: - crate::config::DEFAULT_TUNNEL_STREAM_INITIAL_WINDOW_BYTES, - tunnel_drain_deadline_ms: crate::config::DEFAULT_TUNNEL_DRAIN_DEADLINE_MS, - tunnel_connect_timeout_ms: 2_000, - tunnel_ipv4_only: false, - tunnel_ipv6_only: false, - tunnel_tcp_keepalive_secs: 30, - tunnel_tcp_nodelay: true, - tunnel_stale_timeout_ms: 5_000, - tunnel_connections: Some(1), - tunnel_connections_max: Some(1), - tunnel_scale_check_interval_ms: 1_000, - tunnel_scale_up_threshold_percent: 70, - tunnel_scale_down_threshold_percent: 35, - tunnel_scale_down_grace_secs: 15, - } - } - - fn ensure_rustls_provider() { - static INIT: Once = Once::new(); - INIT.call_once(|| { - let _ = rustls::crypto::ring::default_provider().install_default(); - }); - } } diff --git a/crates/aether-testing/integration/Cargo.toml b/crates/aether-testing/integration/Cargo.toml index 1bc86bb8b..b0d164e8a 100644 --- a/crates/aether-testing/integration/Cargo.toml +++ b/crates/aether-testing/integration/Cargo.toml @@ -15,11 +15,14 @@ aether-data-contracts.workspace = true aether-gateway = { workspace = true, features = ["testkit"] } aether-runtime.workspace = true aether-runtime-state.workspace = true +aether-tunnel.workspace = true aether-testkit = { workspace = true, features = ["gateway", "postgres"] } +arc-swap = "1" axum.workspace = true futures-util.workspace = true http.workspace = true reqwest.workspace = true +rustls.workspace = true serde.workspace = true serde_json.workspace = true sha2.workspace = true diff --git a/crates/aether-testing/integration/tests/tunnel_runtime_e2e.rs b/crates/aether-testing/integration/tests/tunnel_runtime_e2e.rs new file mode 100644 index 000000000..baca13281 --- /dev/null +++ b/crates/aether-testing/integration/tests/tunnel_runtime_e2e.rs @@ -0,0 +1,527 @@ +//! Gateway-backed tunnel end-to-end regressions. +//! +//! 这两个用例需要真实 Gateway 路由和 Tunnel 进程状态,因此放在独立 +//! integration package,避免 Workspace Rest 的普通目标编译 Gateway。 + +use std::future::Future; +use std::pin::Pin; +use std::sync::atomic::AtomicU64; +use std::sync::{Arc, Once}; +use std::task::{Context, Poll}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use aether_contracts::tunnel::{ + sign_tunnel_relay_request, tunnel_relay_payload_digest, TUNNEL_RELAY_AUTH_NONCE_HEADER, + TUNNEL_RELAY_AUTH_PAYLOAD_HEADER, TUNNEL_RELAY_AUTH_SENDER_HEADER, + TUNNEL_RELAY_AUTH_SIGNATURE_HEADER, TUNNEL_RELAY_AUTH_TIMESTAMP_HEADER, + TUNNEL_RELAY_OWNER_INSTANCE_HEADER, +}; +use aether_gateway::{build_router_with_state, AppState as GatewayAppState}; +use aether_tunnel::config::Config; +use aether_tunnel::registration::client::AetherClient; +use aether_tunnel::runtime::DynamicConfig; +use aether_tunnel::state::{ + AppState as TunnelAppState, ServerContext, TunnelMetrics, TunnelRequestMetrics, +}; +use aether_tunnel::target_filter::DnsCache; +use aether_tunnel::tunnel::protocol; +use aether_tunnel::tunnel::run; +use aether_tunnel::upstream_client; +use arc_swap::ArcSwap; +use axum::Router; +use reqwest::StatusCode; +use tokio::sync::watch; + +struct SessionTask(tokio::task::JoinHandle); + +impl SessionTask { + fn new(handle: tokio::task::JoinHandle) -> Self { + Self(handle) + } +} + +impl Future for SessionTask { + type Output = Result; + + fn poll(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll { + Pin::new(&mut self.0).poll(context) + } +} + +impl Drop for SessionTask { + fn drop(&mut self) { + self.0.abort(); + } +} + +#[tokio::test] +async fn tunnel_reconnects_after_gateway_restart() { + ensure_rustls_provider(); + + let gateway_port = reserve_local_port().expect("gateway port should reserve"); + let gateway_base_url = format!("http://127.0.0.1:{gateway_port}"); + let (gateway_state, mut gateway_handle) = start_gateway_on_port(gateway_port) + .await + .expect("gateway should start"); + + let mut tunnel_config = sample_config(&gateway_base_url); + tunnel_config.tunnel_security = aether_tunnel::config::TunnelSecurity::NonTlsRequired; + tunnel_config.tunnel_encryption_key = + Some("BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=".to_string()); + let state = sample_state(tunnel_config); + let server = sample_server(&state, "node-recovery"); + let (shutdown_tx, shutdown_rx) = watch::channel(false); + let tunnel_task = tokio::spawn({ + let state = Arc::clone(&state); + let server = Arc::clone(&server); + let (_drain_tx, drain_rx) = watch::channel(false); + async move { + run(&state, &server, 0, shutdown_rx, drain_rx).await; + } + }); + + wait_until_relay_status( + &gateway_base_url, + "node-recovery", + StatusCode::GATEWAY_TIMEOUT, + ) + .await; + + gateway_handle.abort(); + let _ = (&mut gateway_handle).await; + assert_eq!(gateway_state.force_close_all_tunnel_proxies(), 1); + + let (_restarted_gateway_state, restarted_gateway_handle) = + start_gateway_on_port_retry(gateway_port) + .await + .expect("gateway should restart on fixed port"); + gateway_handle = restarted_gateway_handle; + + wait_until_relay_status( + &gateway_base_url, + "node-recovery", + StatusCode::GATEWAY_TIMEOUT, + ) + .await; + + assert!(server.tunnel_metrics.snapshot().connect_successes >= 2); + let _ = shutdown_tx.send(true); + tokio::time::timeout(Duration::from_secs(5), tunnel_task) + .await + .expect("tunnel task should stop") + .expect("tunnel task should join"); + gateway_handle.abort(); +} + +async fn wait_until_relay_status(gateway_base_url: &str, node_id: &str, expected: StatusCode) { + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let mut last_observed = None::; + loop { + if let Some((status, body)) = probe_relay_status(gateway_base_url, node_id).await { + last_observed = Some(format!("{status} body={body}")); + if status == expected { + return; + } + } + assert!( + tokio::time::Instant::now() < deadline, + "relay status did not become {expected} within timeout; last={:?}", + last_observed + ); + tokio::time::sleep(Duration::from_millis(25)).await; + } +} + +async fn probe_relay_status(gateway_base_url: &str, node_id: &str) -> Option<(StatusCode, String)> { + let response = relay_response(gateway_base_url, node_id, relay_probe_envelope()).await?; + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + Some((status, body)) +} + +async fn relay_response( + gateway_base_url: &str, + node_id: &str, + payload: Vec, +) -> Option { + let timestamp = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("test clock should be after epoch") + .as_secs(); + let nonce = uuid::Uuid::new_v4().simple().to_string(); + let digest = tunnel_relay_payload_digest(&payload, &[]); + let signature = sign_tunnel_relay_request( + b"tunnel-reconnect-test-secret-at-least-32-bytes", + "tunnel-reconnect-test-client", + "tunnel-reconnect-test-gateway", + node_id, + "", + false, + timestamp, + &nonce, + &digest, + ); + reqwest::Client::new() + .post(format!( + "{gateway_base_url}/api/internal/tunnel/relay/{node_id}" + )) + .header("content-type", "application/octet-stream") + .header( + TUNNEL_RELAY_AUTH_SENDER_HEADER, + "tunnel-reconnect-test-client", + ) + .header( + TUNNEL_RELAY_OWNER_INSTANCE_HEADER, + "tunnel-reconnect-test-gateway", + ) + .header(TUNNEL_RELAY_AUTH_TIMESTAMP_HEADER, timestamp) + .header(TUNNEL_RELAY_AUTH_NONCE_HEADER, nonce) + .header( + TUNNEL_RELAY_AUTH_PAYLOAD_HEADER, + digest.encode_header_value(), + ) + .header(TUNNEL_RELAY_AUTH_SIGNATURE_HEADER, signature) + .body(payload) + .send() + .await + .ok() +} + +fn relay_probe_envelope() -> Vec { + let meta = protocol::RequestMeta { + provider_id: None, + endpoint_id: None, + key_id: None, + method: "GET".to_string(), + url: "http://127.0.0.1:80/blocked".to_string(), + headers: std::collections::HashMap::new(), + stream: false, + request_timeout_ms: None, + stream_first_byte_timeout_ms: None, + timeout: 5, + follow_redirects: None, + http1_only: false, + transport_profile: None, + }; + let meta_json = + serde_json::to_vec(&meta).expect("tunnel relay probe metadata should serialize"); + let mut envelope = Vec::with_capacity(4 + meta_json.len()); + envelope.extend_from_slice(&(meta_json.len() as u32).to_be_bytes()); + envelope.extend_from_slice(&meta_json); + envelope +} + +async fn start_gateway_on_port( + port: u16, +) -> Result<(GatewayAppState, tokio::task::JoinHandle<()>), std::io::Error> { + // The embedded gateway now fails closed when relay authentication is + // not configured. Keep this integration fixture explicitly authenticated. + static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + let state = { + let _guard = ENV_LOCK.lock().unwrap(); + let previous_secret = std::env::var_os("AETHER_TUNNEL_RELAY_AUTH_SECRET"); + let previous_instance = std::env::var_os("AETHER_GATEWAY_INSTANCE_ID"); + std::env::set_var( + "AETHER_TUNNEL_RELAY_AUTH_SECRET", + "tunnel-reconnect-test-secret-at-least-32-bytes", + ); + std::env::set_var( + "AETHER_GATEWAY_INSTANCE_ID", + "tunnel-reconnect-test-gateway", + ); + let mut state = GatewayAppState::new().expect("gateway test state should build"); + aether_gateway::configure_test_tunnel_security( + &mut state, + "node-recovery", + "test-generation-1", + "BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=", + ); + restore_test_env("AETHER_TUNNEL_RELAY_AUTH_SECRET", previous_secret); + restore_test_env("AETHER_GATEWAY_INSTANCE_ID", previous_instance); + state + }; + let router = build_router_with_state(state.clone()); + let handle = spawn_router_on_port(port, router).await?; + Ok((state, handle)) +} + +#[tokio::test] +async fn negotiated_small_window_streams_large_responses_and_cancels_idle_upstream() { + use axum::body::{Body, Bytes}; + use axum::routing::get; + use futures_util::StreamExt; + + ensure_rustls_provider(); + let upstream_port = reserve_local_port().unwrap(); + let upstream = Router::new() + .route( + "/large", + get(|| async { Body::from(vec![b'x'; 2 * 1024 * 1024]) }), + ) + .route( + "/idle", + get(|| async { + let first = futures_util::stream::once(async { + Ok::<_, std::io::Error>(Bytes::from_static(b"data: started\n\n")) + }); + ( + [("content-type", "text/event-stream")], + Body::from_stream(first.chain(futures_util::stream::pending())), + ) + }), + ); + let upstream_task = + SessionTask::new(spawn_router_on_port(upstream_port, upstream).await.unwrap()); + let gateway_port = reserve_local_port().unwrap(); + let gateway_url = format!("http://127.0.0.1:{gateway_port}"); + let (_, gateway_task) = start_gateway_on_port(gateway_port).await.unwrap(); + let gateway_task = SessionTask::new(gateway_task); + let mut config = sample_config(&gateway_url); + config.tunnel_security = aether_tunnel::config::TunnelSecurity::NonTlsRequired; + config.tunnel_encryption_key = Some("BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc=".into()); + config.tunnel_stream_initial_window_bytes = 512 * 1024; + config.tunnel_drain_deadline_ms = 100; + config.allow_private_targets = true; + config.allowed_ports.push(upstream_port); + let state = sample_state(config); + let server = sample_server(&state, "node-recovery"); + let (shutdown_tx, shutdown_rx) = watch::channel(false); + let (_drain_tx, drain_rx) = watch::channel(false); + let tunnel_task = SessionTask::new(tokio::spawn({ + let state = Arc::clone(&state); + let server = Arc::clone(&server); + async move { + run(&state, &server, 0, shutdown_rx, drain_rx).await; + } + })); + wait_until_relay_status(&gateway_url, "node-recovery", StatusCode::GATEWAY_TIMEOUT).await; + + let envelope = |path: &str| { + let mut meta: protocol::RequestMeta = + serde_json::from_slice(&relay_probe_envelope()[4..]).unwrap(); + meta.url = format!("http://127.0.0.1:{upstream_port}/{path}"); + meta.stream = true; + meta.timeout = 10; + meta.stream_first_byte_timeout_ms = Some(10_000); + let encoded = serde_json::to_vec(&meta).unwrap(); + let mut result = (encoded.len() as u32).to_be_bytes().to_vec(); + result.extend(encoded); + result + }; + let response = relay_response(&gateway_url, "node-recovery", envelope("large")) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = tokio::time::timeout(Duration::from_secs(10), response.bytes()) + .await + .unwrap() + .unwrap(); + assert_eq!(body.len(), 2 * 1024 * 1024); + assert!(body.iter().all(|byte| *byte == b'x')); + + let mut response = relay_response(&gateway_url, "node-recovery", envelope("idle")) + .await + .unwrap(); + assert_eq!( + response.chunk().await.unwrap().unwrap(), + "data: started\n\n" + ); + drop(response); + tokio::time::timeout(Duration::from_secs(3), async { + while server + .active_connections + .load(std::sync::atomic::Ordering::Acquire) + != 0 + { + tokio::task::yield_now().await; + } + }) + .await + .expect("cancelled SSE must release the upstream handler"); + + let mut response = relay_response(&gateway_url, "node-recovery", envelope("idle")) + .await + .unwrap(); + assert!(response.chunk().await.unwrap().is_some()); + shutdown_tx.send(true).unwrap(); + tokio::time::timeout(Duration::from_secs(3), tunnel_task) + .await + .unwrap() + .unwrap(); + assert_eq!( + server + .active_connections + .load(std::sync::atomic::Ordering::Acquire), + 0 + ); + drop(response); + drop(gateway_task); + drop(upstream_task); +} + +fn restore_test_env(key: &str, value: Option) { + if let Some(value) = value { + std::env::set_var(key, value); + } else { + std::env::remove_var(key); + } +} + +async fn start_gateway_on_port_retry( + port: u16, +) -> Result<(GatewayAppState, tokio::task::JoinHandle<()>), std::io::Error> { + let mut attempts = 0usize; + loop { + match start_gateway_on_port(port).await { + Ok(server) => return Ok(server), + Err(err) => { + attempts += 1; + if attempts >= 20 { + return Err(err); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + } + } +} + +async fn spawn_router_on_port( + port: u16, + app: Router, +) -> Result, std::io::Error> { + let listener = tokio::net::TcpListener::bind(("127.0.0.1", port)).await?; + Ok(tokio::spawn(async move { + axum::serve( + listener, + app.into_make_service_with_connect_info::(), + ) + .await + .expect("gateway test server should run"); + })) +} + +fn reserve_local_port() -> Result { + let listener = std::net::TcpListener::bind("127.0.0.1:0")?; + let port = listener.local_addr()?.port(); + drop(listener); + Ok(port) +} + +fn sample_state(config: Config) -> Arc { + let config = Arc::new(config); + let dns_cache = Arc::new(DnsCache::new(Duration::from_secs(60), 128)); + let upstream_client_pool = + upstream_client::UpstreamClientPool::new(Arc::clone(&config), Arc::clone(&dns_cache)); + Arc::new(TunnelAppState { + config, + dns_cache, + upstream_client_pool, + tunnel_tls_config: Arc::new(aether_tunnel::tunnel::client::build_tls_config()), + resource_monitor: Arc::new(aether_tunnel::hardware::RuntimeResourceMonitor::new()), + stream_gate: None, + distributed_stream_gate: None, + }) +} + +fn sample_server(state: &Arc, node_id: &str) -> Arc { + let config = Arc::clone(&state.config); + Arc::new(ServerContext { + server_label: "gateway-owned-tunnel".to_string(), + aether_url: config.aether_url.clone(), + management_token: config.management_token.clone(), + tunnel_security: config.tunnel_security, + tunnel_encryption_key: config.tunnel_encryption_key.clone(), + node_name: config.node_name.clone(), + node_id: Arc::new(std::sync::RwLock::new(node_id.to_string())), + tunnel_generation: "test-generation-1".to_string(), + aether_client: Arc::new(AetherClient::new( + &config, + &config.aether_url, + &config.management_token, + )), + dynamic: Arc::new(ArcSwap::from_pointee(DynamicConfig::from_config(&config))), + active_connections: Arc::new(AtomicU64::new(0)), + metrics: Arc::new(TunnelRequestMetrics::new()), + tunnel_metrics: Arc::new(TunnelMetrics::new()), + }) +} + +fn sample_config(aether_url: &str) -> Config { + Config { + aether_url: aether_url.to_string(), + management_token: "token".to_string(), + public_ip: None, + node_name: "tunnel-test".to_string(), + tunnel_security: aether_tunnel::config::TunnelSecurity::Off, + tunnel_encryption_key: None, + node_region: None, + heartbeat_interval: 1, + allowed_ports: vec![80, 443], + allow_private_targets: false, + aether_request_timeout_secs: 10, + aether_connect_timeout_secs: 2, + aether_pool_max_idle_per_host: 8, + aether_pool_idle_timeout_secs: 90, + aether_tcp_keepalive_secs: 60, + aether_tcp_nodelay: true, + aether_http2: true, + aether_outbound_proxy_url: None, + aether_retry_max_attempts: 1, + aether_retry_base_delay_ms: 50, + aether_retry_max_delay_ms: 100, + diagnostics_bind: None, + max_concurrent_connections: None, + max_in_flight_streams: None, + distributed_stream_limit: None, + distributed_stream_redis_url: None, + distributed_stream_redis_key_prefix: None, + distributed_stream_lease_ttl_ms: 30_000, + distributed_stream_renew_interval_ms: 10_000, + distributed_stream_command_timeout_ms: 1_000, + dns_cache_ttl_secs: 60, + dns_cache_capacity: 128, + upstream_connect_timeout_secs: 30, + upstream_pool_max_idle_per_host: 4, + upstream_pool_idle_timeout_secs: 60, + upstream_client_pool_capacity: aether_tunnel::config::DEFAULT_UPSTREAM_CLIENT_POOL_CAPACITY, + upstream_tcp_keepalive_secs: 60, + upstream_tcp_nodelay: true, + upstream_proxy_url: None, + upstream_proxy_remote_dns: false, + legacy_redirect_replay_budget_bytes_ignored: None, + emit_proxy_timing_header: true, + log_level: "info".to_string(), + log_destination: aether_tunnel::config::TunnelLogDestinationArg::Stdout, + log_dir: None, + log_rotation: aether_tunnel::config::TunnelLogRotationArg::Daily, + log_retention_days: 7, + log_max_files: 30, + tunnel_reconnect_base_ms: 50, + tunnel_reconnect_max_ms: 250, + tunnel_ping_interval_ms: 1_000, + tunnel_max_streams: Some(8), + tunnel_profile: aether_tunnel::config::TunnelProfileArg::Lite, + tunnel_stream_initial_window_bytes: + aether_tunnel::config::DEFAULT_TUNNEL_STREAM_INITIAL_WINDOW_BYTES, + tunnel_drain_deadline_ms: aether_tunnel::config::DEFAULT_TUNNEL_DRAIN_DEADLINE_MS, + tunnel_connect_timeout_ms: 2_000, + tunnel_ipv4_only: false, + tunnel_ipv6_only: false, + tunnel_tcp_keepalive_secs: 30, + tunnel_tcp_nodelay: true, + tunnel_stale_timeout_ms: 5_000, + tunnel_connections: Some(1), + tunnel_connections_max: Some(1), + tunnel_scale_check_interval_ms: 1_000, + tunnel_scale_up_threshold_percent: 70, + tunnel_scale_down_threshold_percent: 35, + tunnel_scale_down_grace_secs: 15, + } +} + +fn ensure_rustls_provider() { + static INIT: Once = Once::new(); + INIT.call_once(|| { + let _ = rustls::crypto::ring::default_provider().install_default(); + }); +} diff --git a/docs/operations/gateway-ci-timeout-reduction-plan.md b/docs/operations/gateway-ci-timeout-reduction-plan.md new file mode 100644 index 000000000..18be03ccf --- /dev/null +++ b/docs/operations/gateway-ci-timeout-reduction-plan.md @@ -0,0 +1,262 @@ +# Test (Gateway) CI 耗时减负方案 + +- 文档日期:2026-09-23(Asia/Shanghai) +- 目标 job:`.github/workflows/rust-ci.yml` 中的 `test_gateway`(显示名 `Test (Gateway)`) +- 优化目标:**缩短 CI 耗时**,不减少安全/回归断言 +- 关联审计:`docs/operations/system-slimming-audit-2026-09-08.md` +- **第一批(A1+A2+A4)已实施**,待 CI 前后对照确认分钟数。 + +本文只新增方案文档,不修改业务代码、测试或工作流。文中收益均为基于历史日志与源码结构的**估计值**,每批改动落地后须用同一提交做前后对照实测。 + +--- + +## 一、结论 + +`Test (Gateway)` 的耗时大头是**巨型单一 lib 测试目标的编译**(约 60%),其次是 **5300+ 测试的执行**(约 40%)。单纯减少“测试项数”对总分钟数帮助有限;应优先砍掉**重复编译指纹浪费**和**少数重型夹具**。 + +- 单 job 实测约 **11~17 分钟**,是普通 Rust CI 的关键路径。 +- lib 测试约 5300 项、bins 约 78 项,全部 `#[cfg(test)]` 代码编进**同一个** rustc 调用,单进程峰值约 **7GB**,曾出现 OOM。 +- 执行段历史波动 **214s~390s**;其中 4 个慢用例合计约 **54 秒**(约占快样本执行时间的 25%)。 +- 仓库已有瘦身审计的立场不变:**减编译与运行重量,不先减安全断言**。 + +### 不做清单(铁律) + +- 不删除 OAuth、配额、安全头、并发门禁、备份密码学、PII/断开结算等断言。 +- 不下调 `PBKDF2_ITERATIONS`(`crates/aether-crypto/src/python_fernet.rs`)。 +- 不把关键测试标 `ignore`,不跨测试共享可变 `AppState`。 +- 暂不做 nextest 多 runner 分片(每个 runner 各自重编 7GB 巨型目标会净亏;须先复用同一次构建产物再评估)。 + +--- + +## 二、现状基线 + +### 2.1 Job 做什么 + +`.github/workflows/rust-ci.yml` 中 `test_gateway` 关键步骤: + +| 步骤 | 行号 | 命令 / 配置 | +| --- | --- | --- | +| Rust cache | 247-251 | `shared-key: rust-ci-${{ runner.os }}`(与 Data/Rest/Integration 等共用) | +| Setup mold | 256-257 | **仅此 job** 安装 mold | +| Test lib | 265-271 | `cargo nextest run -p aether-gateway --lib` | +| Test bins | 273-279 | `cargo nextest run -p aether-gateway --bins` | + +两个测试 step 的环境变量:`RUSTC_WRAPPER=sccache`、`RUST_MIN_STACK=16777216`、`RUSTFLAGS: "-C link-arg=-fuse-ld=mold"`(**RUSTFLAGS 只在 step 级,未提到 job 级**)。 + +Workflow 级(72-76 行):`CARGO_INCREMENTAL=0`、`CARGO_PROFILE_DEV_DEBUG=0`、`CARGO_PROFILE_TEST_DEBUG=0`。 + +**未发现** `.config/nextest.toml` 或任何 nextest 配置文件;CI 使用 nextest 默认并行度。 + +**未执行**:`apps/aether-gateway/tests/admin_unsigned_identity_headers.rs`(1 个安全集成用例)——两条 nextest 命令只覆盖 `--lib` 与 `--bins`。 + +### 2.2 耗时拆分(历史日志) + +来源:`docs/operations/system-slimming-audit-2026-09-08.md`。 + +| 运行 | Gateway lib 编译 | Gateway lib 执行 | Test lib 步骤合计 | +| --- | --- | --- | --- | +| `34174603131` | 4:57(297s) | **213.807s**(5,139 项) | 514s | +| `34153166516` | 6:22(382s) | **390.039s** | — | + +同轮 bins:编译约 2:27~3:05,执行仅 **0.27~0.33s**(78 项)。 + +- 编译 : 执行 ≈ **57:43**(快样本)~ **49:51**(慢样本)。 +- 合并 `--lib --bins` 为一条命令**不能**省掉普通库与 `cfg(test)` 测试库的两种构建。 + +### 2.3 测试规模(静态统计) + +范围:`apps/aether-gateway`。 + +| 分区 | 约 `#[test]` / `#[tokio::test]` 数量 | +| --- | --- | +| `src/handlers` 内联 | 1224 | +| `src/tests/` 测试树(含 control 727、frontdoor 217、architecture 208 等) | ~1358 | +| `src/execution_runtime` 内联 | ~591 | +| `src/control` 内联 | 459 | +| `src/ai_serving` 内联 | 302 | +| `src/main.rs` + `src/bin/*`(bins 目标) | ~77 | +| 其余分散内联 | ~1500+ | +| **合计** | **约 5380~5400** | + +与 CI 实测(lib 5139 + bins 78)同量级;差异来自 cfg 门控与统计口径。 + +架构守卫:已迁至 `apps/aether-gateway/tests/architecture/`(入口 `tests/architecture_guard.rs`),共 13 文件、约 **14,925 行、208 个 `#[test]`**,断言全部是 `fs::read_to_string` + 字符串/`Cargo.toml` 规则检查,**不依赖 gateway 私有类型**。CI 由 `Test integration targets` 步骤执行(`cargo nextest run -p aether-gateway --tests`)。 + +### 2.4 编译为何是单点瓶颈 + +- `apps/aether-gateway/src/lib.rs:235-236` 把整棵 `tests/` 树挂进**同一个** lib test binary。 +- `src/tests/` 约 130 个文件、17 万行;连同业务代码,单 rustc 调用编译约 53 万行。 +- 历史记录(`openai-responses-websocket-plan.md`):单进程 rustc 峰值约 **7GB**,极端时 RSS 8.9GB 并触发 `oom_kill`;`-j` / `--test-threads` **不影响**该单进程峰值。 +- 对照实验:临时裁掉无关测试模块后,编译可降至约 3 分钟内且零 OOM——证明“编译面”而非“并行度”是主因。 +- 本地 `RUSTC_WRAPPER=sccache` 命中率曾低至约 1.5%;CI 样本 sccache 命中约 95% 仍要数分钟——剩余问题在 **crate-type 不可缓存调用与 RUSTFLAGS 指纹不一致**,不是“再加一层缓存”。 + +### 2.5 执行段慢点(已实测或源码可证) + +| 优先级 | 模式 | 位置 | 估计可省 | +| --- | --- | --- | --- | +| 1 | 池调度超大夹具(1700/2048 账号,每 key 真实 `seal_provider_catalog_key_api_key`) | `src/dispatch/pool_scheduler.rs`(`large_pool_fixture` 约 :5036;慢用例约 :4048、:3967、:4641) | **30~50s**(Top4 中 3 项) | +| 2 | 备份 v1 兼容 + 17 个 historical 密钥各走 10 万次 PBKDF2 | `src/backup/executor.rs:1222` 一带;`python_fernet.rs` | **10~15s**(与上表 14.3s 重叠) | +| 3 | 每测试独立进程重复付 DEVELOPMENT_ENCRYPTION_KEY 的 PBKDF2(nextest 无法跨用例共享缓存) | 全包大量 `encrypt_python_fernet_plaintext(DEVELOPMENT_ENCRYPTION_KEY, …)` | **10~25s**(非密码学夹具改直接密钥后) | +| 4 | 真实 `start_server` 起服极多(tests 树约 1794 次)+ `AppState::new()`(tests 树约 816 次) | `src/tests/mod.rs:33-47` 等 | **15~40s**(改 oneshot 仅限纯路由断言) | +| 5 | 固定 `sleep(100ms)` 扇出、个别 mock 内 5s/30s sleep | `tests/ai_execute/**`、`stream_pump.rs` 等 | **5~20s** | +| 6 | 架构守卫重复全树扫描(迁出后不再计入 lib) | `tests/architecture/*` | 执行 **~3-10s** + 编译见阶段 A | + +**执行优化硬顶**:即使执行砍掉一半,关键路径大约只省 **40~80 秒**;样本间 214s vs 390s 的抖动本身就 ±176s。因此**编译侧(阶段 A/C)才是分钟级收益来源**。 + +--- + +## 三、分阶段方案 + +### 阶段 A:低风险、优先实施(预计单 job 省 1~3 分钟) + +#### A1. 架构守卫迁出独立测试目标【已实施】 + +- **改动**:`src/tests/architecture/**`(208 项 / 约 1.5 万行)已迁至 `apps/aether-gateway/tests/architecture/`,入口 `apps/aether-gateway/tests/architecture_guard.rs`;已从 `src/tests/mod.rs` 移除 `mod architecture;`。 +- **为什么省**:从巨型 lib `cfg(test)` 编译单元削掉约 15k 行,压低 rustc 7GB 峰值,降低 OOM 与编译墙钟;对照实验表明裁测试面可明显缩短编译。 +- **收益估计**:编译 **-20~60s**,执行 **-数秒**;OOM 风险显著下降。 +- **风险**:低。断言逻辑一字不改,仅更换宿主;helper 已改为 `pub(crate)`。 +- **CI**:`test_gateway` 新增 `Test integration targets`:`cargo nextest run -p aether-gateway --tests`(同时覆盖 A4 的 `admin_unsigned_identity_headers`)。 +- **验收**:架构 208 项 + 安全 1 项在新 target 全绿(本地实测 209 passed);lib 测试数减少约 208;lib 编译时间与 rustc 峰值待 CI 对照。 + +#### A2. 统一构建指纹(RUSTFLAGS + 工具链提到 job 级)【已实施】 + +- **改动**: + 1. `test_gateway` 的 mold `RUSTFLAGS`、`RUST_MIN_STACK`、`RUSTC_WRAPPER` 已上移至 **job 级 `env`**,lib/bins/integration 三步共用同一指纹。 + 2. `rust-ci.yml` 中**全部** `Install Rust toolchain` 步骤已钉住 `toolchain: 1.95.0`(与 `rust-toolchain.toml`、fmt/clippy 一致),消除浮动 stable 漂移。 + 3. `test_gateway` 的 `shared-key` 改为 `rust-ci-gateway-test-${{ runner.os }}`,避免 mold 指纹与无 mold 的 job 互相污染共享缓存(一次冷缓存成本)。 +- **为什么省**:原先 step 级 RUSTFLAGS + 浮动 stable + 十余个 job 共用同一 cache key,造成“看似命中、实际重编”,放大 4:57 vs 6:22 的波动;stable 漂移还会触发偶发全量重编。 +- **收益估计**:热缓存编译 **-1~2 分钟波动收窄**;避免偶发 **-5~10 分钟** 尖峰;gateway 独立 cache key 后与 Rest/Data 不再争抢/覆盖。 +- **风险**:低。mold 本就在用;独立 key 首次为冷缓存。 +- **验收**:连续多次 run 的 lib 编译时间方差下降;全 workflow 无未钉 toolchain 步骤。 + +#### A3.(可选,紧随 A2)按用途拆分 rust-cache key + +- **改动**:Gateway **test** 与 lint/check 类 job 不再共用同一 `shared-key`;或评估 `cache-workspace-crates: true`。 +- **收益估计**:编译 **-30~90s**(估计)。 +- **风险**:中。缓存体积上升;勿为每个细碎目标无限新建 key,勿直接缓存完整巨型 `target/`。 + +#### A4. 补跑漏掉的安全集成测试【已实施】 + +- **改动**:`test_gateway` 增加 `Test integration targets`:`cargo nextest run -p aether-gateway --tests`,覆盖 `apps/aether-gateway/tests/` 下全部目标(`architecture_guard` + `admin_unsigned_identity_headers`)。 +- **收益**:时长 **+10~30s**,换取已确认的覆盖缺口(安全优先,与 A1 同 PR)。 +- **风险**:低。该测试走公开 API。本地已实测 1/1 passed。 + +### 阶段 B:只动测试代码(预计执行段省 40~70 秒) + +#### B1. 池调度大夹具轻量化 + +- **改动**:`large_pool_fixture` 对非“规模/扫描预算语义”的用例,改为轻量 repository/credential fixture 或预构造行;**保留**: + - 扫描预算、跳过计数、分页/游标语义断言; + - **至少 1~2 条** 1700/2048 规模边界用例(可保留真实加密封装)。 +- **收益估计**:执行 **-30~50s**。 +- **风险**:中。不得把池规模缩到失去原回归条件;不得删断言。 + +#### B2. PBKDF2 夹具去重 + +- **改动**:对**不测密码学语义**的夹具,改用直接 32 字节 base64 密钥(`decode_direct_fernet_key` 路径)或预制密文,绕开每进程 10 万次迭代。 +- **必须保留**:备份 historical 密钥兼容、生产强度派生、显式 PBKDF2 行为测试(如 `python_fernet` 相关用例)。 +- **收益估计**:执行 **-10~25s**。 +- **风险**:中。**绝不降低迭代次数**。 + +#### B3.(可选,后置)起服与 helper 瘦身 + +- 纯鉴权/路由断言:`start_server` → `tower` oneshot(已有 `send_request`);**涉及超时、连接、完整中间件链的用例保留真实 server**。 +- 合并 8+ 份同构 `run_*_test` 大栈 helper 为单一 helper;OAuth/keys/quota 同构用例可数据驱动,**表内逐行保留断言**。 +- 固定 `sleep` 改为 channel/notify 事件驱动,**不删除等待本身**(防 flaky)。 +- **收益估计**:执行 **-15~40s**;改动面大于 B1/B2,单独 PR。 + +#### B4. nextest 稳定性配置(非直接减分钟) + +- 新增 `.config/nextest.toml`:显式 `test-threads`(对齐 runner vCPU,避免规格漂移)、`slow-timeout`(防单测卡死拖满 job)。 +- **不要**为提速下调 `RUST_MIN_STACK`(16MB 为深栈/管理面用例正确性所需)。 + +### 阶段 C:流水线级(主要缩短累计 runner,间接稳定 Gateway 缓存) + +| # | 改动 | 收益 | 风险 | +| --- | --- | --- | --- | +| C1 | **解除 Tunnel → Gateway dev-dependency**(`apps/aether-tunnel/Cargo.toml:48-49`;端到端场景迁到 integration package) | Rest/Clippy 不再重复编 Gateway → 流水线累计约 **-8~10 分钟 runner**;缓解共享缓存污染 | 中:须迁移 `src/tunnel/mod.rs` 中依赖 `AppState`/`build_router_with_state` 的用例并比对测试清单 | +| C2 | **路径过滤分层**(`rust-ci.yml` push/PR paths):README、安装脚本、Compose、部分 `tests/*.sh` 不再触发全量 Rust 编译 | 非 Rust 变更 **整段跳过 Test (Gateway)**;须保留稳定 gate 防止 required check 永久 pending;补 `rust-toolchain.toml`、`.cargo/**` 触发项 | 中 | +| C3 | Nightly 空 doctest、重复 adapter/feature job 治理(见既有瘦身审计) | 流水线约 **-4 分钟**,**不在**本 job 关键路径 | 低-中 | + +--- + +## 四、实施顺序与验收 + +| 批次 | 内容 | 预期(估计) | 验收标准 | +| --- | --- | --- | --- | +| **第一批 PR** | A1 架构守卫迁出 + A2 指纹统一 + A4 补安全集成测试【已实施】 | 单 job **-1~2 分钟** + 覆盖补齐 | 架构 208 + 安全 1 本地 209 全绿;CI 对照编译时间下降 | +| **第二批 PR** | B1 池调度夹具 + B2 PBKDF2 去重 | 执行 **-40~70s** | 断言集合不减;4 个历史慢用例计时明显下降;密码学用例仍为生产强度 | +| **第三批 PR** | A3 拆缓存 key + C1 Tunnel 解耦 + C2 路径过滤 | 缓存更稳 + 流水线累计大幅下降 | Rest 依赖闭包不再含 Gateway;无关路径 PR 不再拉起全量编译;Gateway 编译方差下降 | +| **可选** | B3 起服/helper、B4 nextest.toml | 再 **-15~40s** + 稳定性 | 无新增 flaky;慢测试超时有告警 | + +### 对照方法(每批必做) + +1. 固定同一 SHA、runner 规格、工具链与 feature 集合。 +2. 记录:`Test lib` / `Test bins` 的**编译墙钟**、**执行墙钟**、测试总数、失败数。 +3. 记录 rustc 峰值内存(如有)与 sccache 命中率。 +4. 区分**冷缓存 / 热缓存**,区分 **job 耗时 / 流水线总耗时**。 +5. **测试总数只允许**因 A1 迁移而在 lib 与新 target 之间搬家;禁止静默减少断言。 + +### 成功指标 + +- **首要**:普通 PR 上 `Test (Gateway)` 墙钟时间下降且结果稳定(方差收窄)。 +- **次要**:全流水线累计 runner 分钟下降(阶段 C)。 +- **禁止**把“删除测试数量”当作成功标准。 + +--- + +## 五、可重复的只读核查命令 + +```sh +# 测试属性数量(lib 树 / bins) +rg -c '#\[(tokio::)?test\]' apps/aether-gateway/src -g '*.rs' | awk -F: '{s+=$2} END {print s}' +rg -c '#\[(tokio::)?test\]' apps/aether-gateway/src/main.rs apps/aether-gateway/src/bin -g '*.rs' + +# 架构守卫规模(迁移后) +rg -c '#\[(tokio::)?test\]' apps/aether-gateway/tests/architecture -g '*.rs' +wc -l apps/aether-gateway/tests/architecture/*.rs + +# CI 集成测试步骤与指纹 +rg -n 'Test integration targets|rust-ci-gateway-test|toolchain: 1.95.0|RUSTFLAGS' .github/workflows/rust-ci.yml + +# 是否存在 nextest 配置 +ls .config/nextest.toml 2>/dev/null || echo 'no nextest.toml' + +# Tunnel 反向依赖 +rg -n 'aether-gateway' apps/aether-tunnel/Cargo.toml + +# CI 中 mold/RUSTFLAGS/缓存 key +rg -n 'mold|RUSTFLAGS|shared-key|nextest run -p aether-gateway' .github/workflows/rust-ci.yml +``` + +历史耗时与慢用例计时以 `docs/operations/system-slimming-audit-2026-09-08.md` 及对应 GitHub Actions 运行 ID 为准;临时 API JSON 不入库。 + +--- + +## 六、风险与回滚 + +| 风险 | 缓解 | +| --- | --- | +| 迁出架构守卫后漏挂模块 | 迁移前后对比 208 项清单;CI 显式跑新 target | +| RUSTFLAGS 上移后某 job 链接失败 | 先在 `test_gateway` job 级验证,再推广到其它 job;保留 step 级回滚 diff | +| 夹具轻量化导致规模回归失效 | 强制保留 ≥1 条大规模边界用例;PR 中 diff 审查断言 | +| 路径过滤导致 required check 永久 pending | 增加始终运行的轻量 gate job | +| Tunnel 解耦丢失端到端场景 | 迁移前后测试清单比对;场景迁入 integration job | + +回滚单位按 **PR 批次**:每批独立可 revert,不把 A/B/C 混在同一提交。 + +--- + +## 七、附录:明确不能减的测试(摘录) + +| 类别 | 主要位置(约) | 说明 | +| --- | --- | --- | +| OAuth 导入/刷新/吊销 | `src/tests/control/admin/oauth.rs` 等 | 账户安全核心 | +| 配额与失效删 key | `src/tests/control/admin/endpoints/quota.rs` 等 | 计费与授权正确性 | +| 安全头 / 管理面访问 | `security.rs`、`health_access.rs`、`operational_auth` | 未签名身份头等 | +| 并发门禁 | `src/tests/concurrency.rs` | 过载拒绝与独立准入 | +| 备份历史兼容与密码学 | `src/backup/executor.rs` | 保留生产强度 PBKDF2 | +| 池调度扫描预算语义 | `src/dispatch/pool_scheduler.rs` | 可改夹具实现,不可删语义断言 | +| AI 断开结算 / PII | `src/tests/ai_execute/...` | 计费完整性与隐私 | + +> 与 `system-slimming-audit-2026-09-08.md` 一致:首要结果是 **PR 更快得到正确反馈**,不是测试条数变少。