diff --git a/.github/workflows/rust-ci.yml b/.github/workflows/rust-ci.yml index 5b96b3813..3f4149def 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,6 +245,8 @@ 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 指纹不一致。 @@ -289,6 +302,8 @@ jobs: test_data: name: Test (Data) + needs: changes + if: ${{ needs.changes.outputs.rust == 'true' }} runs-on: ubuntu-latest steps: - uses: actions/checkout@fbc6f3992d24b796d5a048ff273f7fcc4a7b6c09 # v5 @@ -332,6 +347,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 @@ -371,6 +388,8 @@ 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 @@ -410,6 +429,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 @@ -451,6 +472,8 @@ 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 @@ -489,6 +512,7 @@ jobs: name: Test runs-on: ubuntu-latest needs: + - changes - test_gateway - test_data - check_data_features @@ -499,6 +523,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" ] || \ @@ -511,6 +539,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: @@ -600,11 +630,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 @@ -614,6 +649,7 @@ jobs: name: check runs-on: ubuntu-latest needs: + - changes - fmt - clippy - test @@ -623,11 +659,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 e475509d2..3f8eee4ad 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", @@ -671,7 +674,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-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(); + }); +}