fix(http-tunnel): reduce round-trip latency under active traffic

This commit is contained in:
Zzz
2026-08-26 22:00:04 +08:00
committed by GitHub
parent ac2af84e7b
commit 4b454b72b9
3 changed files with 233 additions and 11 deletions
+31 -10
View File
@@ -301,9 +301,11 @@ fn effective_connect_timeout_secs(value: u64) -> u64 {
#[cfg(test)]
mod tests {
use super::{script_url, validate_script_url, HttpTunnelManager};
use std::sync::Arc;
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
use tokio::net::{TcpListener, TcpStream};
use tokio::sync::Mutex;
use tokio::sync::{Mutex, Notify};
use tokio::time::{timeout, Duration};
#[test]
fn script_url_preserves_existing_query_and_appends_action() {
@@ -331,16 +333,28 @@ mod tests {
}
#[tokio::test]
async fn manager_forwards_bytes_through_http_script_protocol() {
async fn manager_forwards_sequential_round_trips_and_closes_session() {
let script = MockScript::start().await;
let manager = HttpTunnelManager::new();
let local_port = manager.start_tunnel("test", &script.url, "secret", 5, "mysql.internal", 3306).await.unwrap();
let mut client = TcpStream::connect(("127.0.0.1", local_port)).await.unwrap();
client.write_all(b"ping").await.unwrap();
let mut response = [0_u8; 4];
client.read_exact(&mut response).await.unwrap();
let close_notification = script.closed.notified();
assert_eq!(&response, b"ping");
timeout(Duration::from_secs(10), async {
for index in 0..100_u32 {
let payload = format!("request-{index:03}-{}", "x".repeat((index % 17) as usize));
client.write_all(payload.as_bytes()).await.unwrap();
let mut response = vec![0_u8; payload.len()];
client.read_exact(&mut response).await.unwrap();
assert_eq!(response, payload.as_bytes());
}
})
.await
.expect("100 sequential HTTP tunnel round trips timed out");
drop(client);
timeout(Duration::from_secs(2), close_notification).await.expect("HTTP tunnel session was not closed");
manager.stop_tunnel("test").await;
script.handle.abort();
@@ -348,6 +362,7 @@ mod tests {
struct MockScript {
url: String,
closed: Arc<Notify>,
handle: tokio::task::JoinHandle<()>,
}
@@ -355,26 +370,29 @@ mod tests {
async fn start() -> Self {
let listener = TcpListener::bind(("127.0.0.1", 0)).await.unwrap();
let addr = listener.local_addr().unwrap();
let buffered = std::sync::Arc::new(Mutex::new(Vec::<u8>::new()));
let buffered = Arc::new(Mutex::new(Vec::<u8>::new()));
let closed = Arc::new(Notify::new());
let handle = {
let buffered = buffered.clone();
let closed = closed.clone();
tokio::spawn(async move {
loop {
let Ok((stream, _)) = listener.accept().await else {
break;
};
let buffered = buffered.clone();
let closed = closed.clone();
tokio::spawn(async move {
handle_mock_http_request(stream, buffered).await;
handle_mock_http_request(stream, buffered, closed).await;
});
}
})
};
Self { url: format!("http://{addr}/dbx_tunnel.php"), handle }
Self { url: format!("http://{addr}/dbx_tunnel.php"), closed, handle }
}
}
async fn handle_mock_http_request(stream: TcpStream, buffered: std::sync::Arc<Mutex<Vec<u8>>>) {
async fn handle_mock_http_request(stream: TcpStream, buffered: Arc<Mutex<Vec<u8>>>, closed: Arc<Notify>) {
let mut reader = BufReader::new(stream);
let mut request_line = String::new();
reader.read_line(&mut request_line).await.unwrap();
@@ -407,6 +425,9 @@ mod tests {
} else {
write_http_response(&mut stream, "200 OK", &bytes).await;
}
} else if request_line.contains("dbx_action=close") {
write_http_response(&mut stream, "200 OK", b"OK").await;
closed.notify_one();
} else {
write_http_response(&mut stream, "200 OK", b"OK").await;
}
+38 -1
View File
@@ -21,6 +21,18 @@ $DBX_TUNNEL_DIR = getenv('DBX_TUNNEL_DIR') ?: sys_get_temp_dir() . DIRECTORY_SEP
$DBX_TUNNEL_ALLOWED_HOSTS = array_values(array_filter(array_map('trim', explode(',', getenv('DBX_TUNNEL_ALLOWED_HOSTS') ?: ''))));
$DBX_TUNNEL_MAX_SESSION_SECONDS = max(30, (int) (getenv('DBX_TUNNEL_MAX_SESSION_SECONDS') ?: '3600'));
// Keep sequential protocol exchanges responsive, then restore the original idle cadence.
const DBX_WORKER_ACTIVE_POLL_US = 10000;
const DBX_WORKER_WARM_POLL_US = 50000;
const DBX_WORKER_IDLE_POLL_US = 200000;
const DBX_WORKER_ACTIVE_POLL_COUNT = 100;
const DBX_WORKER_WARM_POLL_COUNT = 20;
// Standalone tests load the worker helpers without dispatching an HTTP request.
if (defined('DBX_TUNNEL_FUNCTIONS_ONLY') && DBX_TUNNEL_FUNCTIONS_ONLY) {
return;
}
if (PHP_SAPI === 'cli' && isset($argv[1]) && $argv[1] === '--dbx-worker') {
run_worker($argv[2] ?? '', $argv[3] ?? '', (int) ($argv[4] ?? 0), (int) ($argv[5] ?? 10));
exit;
@@ -160,6 +172,7 @@ function run_worker(string $dir, string $host, int $port, int $connectTimeout):
stream_set_blocking($socket, false);
$expiresAt = time() + $DBX_TUNNEL_MAX_SESSION_SECONDS;
$lastActivity = time();
$idlePolls = 0;
try {
while (time() < $expiresAt) {
@@ -167,16 +180,18 @@ function run_worker(string $dir, string $host, int $port, int $connectTimeout):
break;
}
$activity = false;
$inbound = drain_chunks($dir, 'in');
if ($inbound !== '') {
write_all($socket, $inbound);
$lastActivity = time();
$activity = true;
}
$read = [$socket];
$write = [];
$except = [];
$ready = @stream_select($read, $write, $except, 0, 200000);
$ready = @stream_select($read, $write, $except, 0, worker_poll_timeout_us($idlePolls));
if ($ready === false) {
write_error($dir, 'Failed to poll target database socket');
break;
@@ -194,9 +209,12 @@ function run_worker(string $dir, string $host, int $port, int $connectTimeout):
} else {
append_chunk($dir, 'out', $data);
$lastActivity = time();
$activity = true;
}
}
$idlePolls = next_worker_idle_poll_count($idlePolls, $activity);
if (time() - $lastActivity > $DBX_TUNNEL_MAX_SESSION_SECONDS) {
break;
}
@@ -209,6 +227,25 @@ function run_worker(string $dir, string $host, int $port, int $connectTimeout):
mark_closed($dir);
}
function worker_poll_timeout_us(int $idlePolls): int
{
if ($idlePolls < DBX_WORKER_ACTIVE_POLL_COUNT) {
return DBX_WORKER_ACTIVE_POLL_US;
}
if ($idlePolls < DBX_WORKER_ACTIVE_POLL_COUNT + DBX_WORKER_WARM_POLL_COUNT) {
return DBX_WORKER_WARM_POLL_US;
}
return DBX_WORKER_IDLE_POLL_US;
}
function next_worker_idle_poll_count(int $idlePolls, bool $activity): int
{
if ($activity) {
return 0;
}
return min($idlePolls + 1, DBX_WORKER_ACTIVE_POLL_COUNT + DBX_WORKER_WARM_POLL_COUNT);
}
function write_all($socket, string $data): void
{
$offset = 0;
+164
View File
@@ -0,0 +1,164 @@
<?php
declare(strict_types=1);
define('DBX_TUNNEL_FUNCTIONS_ONLY', true);
require dirname(__DIR__) . DIRECTORY_SEPARATOR . 'dbx_tunnel.php';
assert_same(10000, worker_poll_timeout_us(0), 'active polling starts at 10ms');
assert_same(10000, worker_poll_timeout_us(99), 'active polling lasts for 100 idle polls');
assert_same(50000, worker_poll_timeout_us(100), 'brief inactivity backs off to 50ms');
assert_same(50000, worker_poll_timeout_us(119), 'warm polling lasts for 20 idle polls');
assert_same(200000, worker_poll_timeout_us(120), 'idle polling returns to 200ms');
assert_same(0, next_worker_idle_poll_count(120, true), 'activity resets the idle counter');
assert_same(10000, worker_poll_timeout_us(next_worker_idle_poll_count(120, true)), 'activity resets polling to 10ms');
assert_same(120, next_worker_idle_poll_count(120, false), 'the idle counter remains capped');
$baseDir = sys_get_temp_dir() . DIRECTORY_SEPARATOR . 'dbx-tunnel-test-' . bin2hex(random_bytes(6));
$sessionDir = $baseDir . DIRECTORY_SEPARATOR . 'session123';
$worker = null;
$pipes = [];
$server = null;
$socket = null;
try {
if (!mkdir($sessionDir, 0700, true) && !is_dir($sessionDir)) {
throw new RuntimeException('Failed to create test session directory');
}
touch($sessionDir . DIRECTORY_SEPARATOR . 'in.queue');
touch($sessionDir . DIRECTORY_SEPARATOR . 'out.queue');
$server = stream_socket_server('tcp://127.0.0.1:0', $errno, $errstr);
if ($server === false) {
throw new RuntimeException('Failed to start echo target: ' . $errstr);
}
$address = stream_socket_get_name($server, false);
$port = (int) substr((string) strrchr($address, ':'), 1);
$script = dirname(__DIR__) . DIRECTORY_SEPARATOR . 'dbx_tunnel.php';
$command = escapeshellarg(PHP_BINARY)
. ' ' . escapeshellarg($script)
. ' --dbx-worker ' . escapeshellarg($sessionDir)
. ' 127.0.0.1 ' . escapeshellarg((string) $port)
. ' 5';
$worker = proc_open($command, [
0 => ['pipe', 'r'],
1 => ['pipe', 'w'],
2 => ['pipe', 'w'],
], $pipes);
if (!is_resource($worker)) {
throw new RuntimeException('Failed to start tunnel worker');
}
fclose($pipes[0]);
stream_set_blocking($pipes[1], false);
stream_set_blocking($pipes[2], false);
$socket = @stream_socket_accept($server, 5);
if ($socket === false) {
throw new RuntimeException('Tunnel worker did not connect to echo target');
}
stream_set_blocking($socket, false);
$roundTrips = 100;
$startedAt = microtime(true);
for ($index = 0; $index < $roundTrips; $index++) {
$payload = pack('N', $index) . hash('sha256', (string) $index, true);
append_chunk($sessionDir, 'in', $payload);
$forwarded = read_exact_with_timeout($socket, strlen($payload), 2.0);
assert_same($payload, $forwarded, 'TCP target receives ordered payload ' . $index);
write_all($socket, $forwarded);
$reply = drain_exact_with_timeout($sessionDir, 'out', strlen($payload), 2.0);
assert_same($payload, $reply, 'tunnel returns ordered payload ' . $index);
}
$elapsed = microtime(true) - $startedAt;
if ($elapsed >= 10.0) {
throw new RuntimeException(sprintf(
'Sequential tunnel round trips took %.3fs; the old fixed wait pattern takes about 20s',
$elapsed
));
}
touch($sessionDir . DIRECTORY_SEPARATOR . 'close');
wait_for_path($sessionDir . DIRECTORY_SEPARATOR . 'closed', 3.0);
printf("ok - %d sequential round trips in %.3fs; ordered payloads and clean shutdown verified\n", $roundTrips, $elapsed);
} finally {
if (is_resource($socket)) {
fclose($socket);
}
if (is_resource($server)) {
fclose($server);
}
if (is_resource($worker)) {
foreach ([1, 2] as $pipeIndex) {
if (isset($pipes[$pipeIndex]) && is_resource($pipes[$pipeIndex])) {
fclose($pipes[$pipeIndex]);
}
}
$status = proc_get_status($worker);
if ($status['running']) {
proc_terminate($worker);
}
proc_close($worker);
}
if (is_dir($baseDir)) {
remove_dir($baseDir);
}
}
function read_exact_with_timeout($socket, int $length, float $timeoutSeconds): string
{
$data = '';
$deadline = microtime(true) + $timeoutSeconds;
while (strlen($data) < $length) {
$chunk = fread($socket, $length - strlen($data));
if ($chunk === false) {
throw new RuntimeException('Failed to read echo target socket');
}
if ($chunk !== '') {
$data .= $chunk;
continue;
}
if (microtime(true) >= $deadline) {
throw new RuntimeException('Timed out reading echo target socket');
}
usleep(1000);
}
return $data;
}
function drain_exact_with_timeout(string $dir, string $name, int $length, float $timeoutSeconds): string
{
$data = '';
$deadline = microtime(true) + $timeoutSeconds;
while (strlen($data) < $length) {
$data .= drain_chunks($dir, $name);
if (strlen($data) >= $length) {
break;
}
if (microtime(true) >= $deadline) {
throw new RuntimeException('Timed out draining tunnel queue');
}
usleep(1000);
}
return $data;
}
function wait_for_path(string $path, float $timeoutSeconds): void
{
$deadline = microtime(true) + $timeoutSeconds;
while (!is_file($path)) {
if (microtime(true) >= $deadline) {
throw new RuntimeException('Timed out waiting for ' . $path);
}
usleep(10000);
}
}
function assert_same($expected, $actual, string $message): void
{
if ($expected !== $actual) {
throw new RuntimeException($message);
}
}