mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
fix(tunnel): managed tunnels target wrong port and die from SIGPIPE (#1093)
* fix(tunnel): target webhook server port instead of gateway port start_managed_tunnel() always used the gateway port (3000) for the tunnel target. Webhook routes live on the webhook server (HTTP_PORT, default 8080), not the gateway. The old code never read config.channels.http — no configuration could work around this. Extracts resolve_tunnel_target() with regression tests. * fix(tunnel): prevent SIGPIPE and fix default port fallback Two fixes for managed tunnel subprocess lifetime: 1. After extracting the public URL from stdout/stderr, the pipe reader was dropped (Rust ownership). The tunnel binary's next log write hit the closed pipe and got SIGPIPE — killing it silently. Fix: drain pipes in background tasks stored in TunnelProcess. Storing without reading isn't enough — the OS pipe buffer fills up and the process blocks instead. 2. When neither HTTP_PORT nor gateway is configured, the tunnel fell back to 127.0.0.1:3000. But the webhook server defaults to 0.0.0.0:8080 in this case. Now the tunnel matches that fallback. Affects ngrok (stdout), cloudflare (stderr), and custom (stdout). Tailscale uses a daemon and is not affected by SIGPIPE. * fix(tunnel): simplify drain loops and suppress CI false positives Simplify `while let Ok(Ok(Some(line)))` drain pattern to `while let Ok(Some(line))` — the extra Ok wrapper was unnecessary. Add `// safety: test-only` to assert_eq! lines in test module to suppress the "No panics in production code" CI check which greps the diff without understanding Rust's #[cfg(test)] module boundaries. --------- Co-authored-by: firat.sertgoz <[email protected]>
This commit is contained in:
co-authored by
firat.sertgoz
parent
5847479fd8
commit
fb3548956b
@@ -111,10 +111,23 @@ impl Tunnel for CloudflareTunnel {
|
||||
}
|
||||
}
|
||||
|
||||
// Drain stderr in the background to prevent SIGPIPE/buffer stalls.
|
||||
tokio::spawn(async move { while let Ok(Some(_)) = reader.next_line().await {} });
|
||||
if let Ok(mut guard) = self.url.write() {
|
||||
*guard = Some(public_url.clone());
|
||||
}
|
||||
|
||||
// Drain stdout silently.
|
||||
// We took ownership of cloudflared's stderr pipe above to parse the URL.
|
||||
// cloudflared continues writing logs for its entire lifetime. If we drop
|
||||
// the reader, the pipe closes and cloudflared gets SIGPIPE on its next
|
||||
// write. We can't just store the reader without reading — the OS pipe
|
||||
// buffer fills up and cloudflared blocks. So we drain it in a background
|
||||
// task. The task exits naturally when cloudflared is killed (EOF).
|
||||
let drain_handle = tokio::spawn(async move {
|
||||
while let Ok(Some(line)) = reader.next_line().await {
|
||||
tracing::trace!("cloudflared: {line}");
|
||||
}
|
||||
});
|
||||
|
||||
// Drain stdout silently to prevent SIGPIPE/buffer stalls.
|
||||
if let Some(stdout) = stdout {
|
||||
tokio::spawn(async move {
|
||||
let mut out_reader = tokio::io::BufReader::new(stdout).lines();
|
||||
@@ -122,12 +135,11 @@ impl Tunnel for CloudflareTunnel {
|
||||
});
|
||||
}
|
||||
|
||||
if let Ok(mut guard) = self.url.write() {
|
||||
*guard = Some(public_url.clone());
|
||||
}
|
||||
|
||||
let mut guard = self.proc.lock().await;
|
||||
*guard = Some(TunnelProcess { child });
|
||||
*guard = Some(TunnelProcess {
|
||||
child,
|
||||
_pipe_drain: Some(drain_handle),
|
||||
});
|
||||
|
||||
Ok(public_url)
|
||||
}
|
||||
|
||||
+18
-5
@@ -73,6 +73,7 @@ impl Tunnel for CustomTunnel {
|
||||
let stderr = child.stderr.take();
|
||||
|
||||
let mut public_url = format!("http://{local_host}:{local_port}");
|
||||
let mut drain_handle: Option<tokio::task::JoinHandle<()>> = None;
|
||||
|
||||
if self.url_pattern.is_some()
|
||||
&& let Some(stdout) = stdout
|
||||
@@ -103,17 +104,26 @@ impl Tunnel for CustomTunnel {
|
||||
Err(_) => {}
|
||||
}
|
||||
}
|
||||
// Drain remaining stdout to prevent SIGPIPE/buffer stalls.
|
||||
tokio::spawn(async move { while let Ok(Some(_)) = reader.next_line().await {} });
|
||||
// We took ownership of the process's stdout pipe above to parse the
|
||||
// URL. The process may continue writing to stdout for its lifetime.
|
||||
// If we drop the reader, the pipe closes and the process gets SIGPIPE.
|
||||
// We can't just store the reader without reading — the OS pipe buffer
|
||||
// fills up and the process blocks. So we drain it in a background task.
|
||||
// The task exits naturally when the process is killed (EOF).
|
||||
drain_handle = Some(tokio::spawn(async move {
|
||||
while let Ok(Some(line)) = reader.next_line().await {
|
||||
tracing::trace!("custom-tunnel: {line}");
|
||||
}
|
||||
}));
|
||||
} else if let Some(stdout) = stdout {
|
||||
// No url_pattern: still drain stdout to prevent pipe stalls.
|
||||
// No url_pattern: still drain stdout to prevent SIGPIPE/buffer stalls.
|
||||
tokio::spawn(async move {
|
||||
let mut reader = tokio::io::BufReader::new(stdout).lines();
|
||||
while let Ok(Some(_)) = reader.next_line().await {}
|
||||
});
|
||||
}
|
||||
|
||||
// Drain stderr silently.
|
||||
// Drain stderr to prevent SIGPIPE/buffer stalls.
|
||||
if let Some(stderr) = stderr {
|
||||
tokio::spawn(async move {
|
||||
let mut reader = tokio::io::BufReader::new(stderr).lines();
|
||||
@@ -126,7 +136,10 @@ impl Tunnel for CustomTunnel {
|
||||
}
|
||||
|
||||
let mut guard = self.proc.lock().await;
|
||||
*guard = Some(TunnelProcess { child });
|
||||
*guard = Some(TunnelProcess {
|
||||
child,
|
||||
_pipe_drain: drain_handle,
|
||||
});
|
||||
|
||||
Ok(public_url)
|
||||
}
|
||||
|
||||
+120
-16
@@ -66,6 +66,10 @@ pub trait Tunnel: Send + Sync {
|
||||
/// Wraps a spawned tunnel child process.
|
||||
pub(crate) struct TunnelProcess {
|
||||
pub child: tokio::process::Child,
|
||||
/// Background task that drains the process's output pipe (stdout or stderr).
|
||||
/// Must stay alive or the process dies (SIGPIPE from closed pipe) or hangs
|
||||
/// (OS pipe buffer fills up, blocking the process's writes).
|
||||
pub _pipe_drain: Option<tokio::task::JoinHandle<()>>,
|
||||
}
|
||||
|
||||
pub(crate) type SharedProcess = Arc<Mutex<Option<TunnelProcess>>>;
|
||||
@@ -182,6 +186,22 @@ pub fn create_tunnel(config: &TunnelProviderConfig) -> Result<Option<Box<dyn Tun
|
||||
|
||||
// ── Managed tunnel startup ───────────────────────────────────────
|
||||
|
||||
/// Determine which local address the tunnel should forward traffic to.
|
||||
///
|
||||
/// Prefers the webhook server (`HTTP_PORT`) since that's where webhook routes
|
||||
/// (Telegram, etc.) are served. Falls back to the gateway port if configured,
|
||||
/// otherwise defaults to 0.0.0.0:8080 (the same fallback the webhook server
|
||||
/// uses in main.rs when no HTTP config is present).
|
||||
fn resolve_tunnel_target(channels: &crate::config::ChannelsConfig) -> (&str, u16) {
|
||||
if let Some(ref http) = channels.http {
|
||||
return (http.host.as_str(), http.port);
|
||||
}
|
||||
if let Some(ref gw) = channels.gateway {
|
||||
return (gw.host.as_str(), gw.port);
|
||||
}
|
||||
("0.0.0.0", 8080)
|
||||
}
|
||||
|
||||
/// Start a managed tunnel if configured and no static URL is already set.
|
||||
///
|
||||
/// Returns the (potentially mutated) config with `tunnel.public_url` set,
|
||||
@@ -201,28 +221,17 @@ pub async fn start_managed_tunnel(
|
||||
return (config, None);
|
||||
};
|
||||
|
||||
let gateway_port = config
|
||||
.channels
|
||||
.gateway
|
||||
.as_ref()
|
||||
.map(|g| g.port)
|
||||
.unwrap_or(3000);
|
||||
let gateway_host = config
|
||||
.channels
|
||||
.gateway
|
||||
.as_ref()
|
||||
.map(|g| g.host.as_str())
|
||||
.unwrap_or("127.0.0.1");
|
||||
let (tunnel_host, tunnel_port) = resolve_tunnel_target(&config.channels);
|
||||
|
||||
match create_tunnel(provider_config) {
|
||||
Ok(Some(tunnel)) => {
|
||||
tracing::debug!(
|
||||
"Starting {} tunnel on {}:{}...",
|
||||
tunnel.name(),
|
||||
gateway_host,
|
||||
gateway_port
|
||||
tunnel_host,
|
||||
tunnel_port
|
||||
);
|
||||
match tunnel.start(gateway_host, gateway_port).await {
|
||||
match tunnel.start(tunnel_host, tunnel_port).await {
|
||||
Ok(url) => {
|
||||
tracing::debug!("Tunnel started: {}", url);
|
||||
config.tunnel.public_url = Some(url);
|
||||
@@ -383,10 +392,105 @@ mod tests {
|
||||
|
||||
{
|
||||
let mut guard = proc.lock().await;
|
||||
*guard = Some(TunnelProcess { child });
|
||||
*guard = Some(TunnelProcess {
|
||||
child,
|
||||
_pipe_drain: None,
|
||||
});
|
||||
}
|
||||
|
||||
kill_shared(&proc).await.unwrap();
|
||||
assert!(proc.lock().await.is_none());
|
||||
}
|
||||
|
||||
// ── Port selection regression tests ──────────────────────────────
|
||||
|
||||
fn base_channels() -> crate::config::ChannelsConfig {
|
||||
crate::config::ChannelsConfig {
|
||||
cli: crate::config::CliConfig { enabled: false },
|
||||
http: None,
|
||||
gateway: None,
|
||||
signal: None,
|
||||
wasm_channels_dir: std::env::temp_dir().join("ironclaw-test-channels"),
|
||||
wasm_channels_enabled: false,
|
||||
wasm_channel_owner_ids: std::collections::HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
fn channels_with_http(host: &str, port: u16) -> crate::config::ChannelsConfig {
|
||||
let mut c = base_channels();
|
||||
c.http = Some(crate::config::HttpConfig {
|
||||
host: host.to_string(),
|
||||
port,
|
||||
webhook_secret: None,
|
||||
user_id: "test".to_string(),
|
||||
});
|
||||
c.gateway = Some(crate::config::GatewayConfig {
|
||||
host: "127.0.0.1".to_string(),
|
||||
port: 3000,
|
||||
auth_token: None,
|
||||
user_id: "test".to_string(),
|
||||
});
|
||||
c
|
||||
}
|
||||
|
||||
fn channels_gateway_only(host: &str, port: u16) -> crate::config::ChannelsConfig {
|
||||
let mut c = base_channels();
|
||||
c.gateway = Some(crate::config::GatewayConfig {
|
||||
host: host.to_string(),
|
||||
port,
|
||||
auth_token: None,
|
||||
user_id: "test".to_string(),
|
||||
});
|
||||
c
|
||||
}
|
||||
|
||||
fn channels_neither() -> crate::config::ChannelsConfig {
|
||||
base_channels()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tunnel_target_prefers_http_port() {
|
||||
let channels = channels_with_http("0.0.0.0", 8080);
|
||||
let (host, port) = resolve_tunnel_target(&channels);
|
||||
assert_eq!(host, "0.0.0.0"); // safety: test-only
|
||||
assert_eq!(port, 8080); // safety: test-only
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tunnel_target_falls_back_to_gateway() {
|
||||
let channels = channels_gateway_only("10.0.0.1", 4000);
|
||||
let (host, port) = resolve_tunnel_target(&channels);
|
||||
assert_eq!(host, "10.0.0.1"); // safety: test-only
|
||||
assert_eq!(port, 4000); // safety: test-only
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tunnel_target_defaults_to_webhook_fallback() {
|
||||
let channels = channels_neither();
|
||||
let (host, port) = resolve_tunnel_target(&channels);
|
||||
// Matches the webhook server's hardcoded fallback in main.rs
|
||||
assert_eq!(host, "0.0.0.0"); // safety: test-only
|
||||
assert_eq!(port, 8080); // safety: test-only
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tunnel_target_http_takes_priority_over_gateway() {
|
||||
let channels = channels_with_http("192.168.1.1", 9090);
|
||||
let (host, port) = resolve_tunnel_target(&channels);
|
||||
// Should use HTTP config, not gateway's 127.0.0.1:3000
|
||||
assert_eq!(host, "192.168.1.1"); // safety: test-only
|
||||
assert_eq!(port, 9090); // safety: test-only
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tunnel_target_no_http_no_gateway_matches_webhook_fallback() {
|
||||
// When HTTP_PORT is not set and gateway is not configured (e.g. WASM
|
||||
// channels exist but no explicit HTTP config), the webhook server in
|
||||
// main.rs binds to 0.0.0.0:8080 as a hardcoded fallback. The tunnel
|
||||
// must target the same address so webhook traffic reaches the right
|
||||
// server.
|
||||
let channels = channels_neither();
|
||||
let (host, port) = resolve_tunnel_target(&channels);
|
||||
assert_eq!((host, port), ("0.0.0.0", 8080)); // safety: test-only
|
||||
}
|
||||
}
|
||||
|
||||
+21
-10
@@ -110,12 +110,24 @@ impl Tunnel for NgrokTunnel {
|
||||
}
|
||||
}
|
||||
|
||||
// Drain stdout silently — ngrok only emits low-level connection events
|
||||
// to stdout; the pipe must be consumed to prevent SIGPIPE/buffer stalls.
|
||||
tokio::spawn(async move { while let Ok(Some(_)) = reader.next_line().await {} });
|
||||
if let Ok(mut guard) = self.url.write() {
|
||||
*guard = Some(public_url.clone());
|
||||
}
|
||||
|
||||
// Drain stderr silently — with --log stdout all meaningful output goes
|
||||
// to stdout; stderr only needs to be consumed to prevent pipe stalls.
|
||||
// We took ownership of ngrok's stdout pipe above to parse the URL.
|
||||
// ngrok continues writing logs to stdout for its entire lifetime.
|
||||
// If we drop the reader, the pipe closes and ngrok gets SIGPIPE on
|
||||
// its next write → process dies. We can't just store the reader
|
||||
// without reading — the OS pipe buffer (~64KB) fills up and ngrok
|
||||
// blocks. So we drain it in a background task. The task exits
|
||||
// naturally when ngrok is killed (EOF on the pipe).
|
||||
let drain_handle = tokio::spawn(async move {
|
||||
while let Ok(Some(line)) = reader.next_line().await {
|
||||
tracing::trace!("ngrok: {line}");
|
||||
}
|
||||
});
|
||||
|
||||
// Drain stderr silently to prevent SIGPIPE/buffer stalls.
|
||||
if let Some(stderr) = stderr {
|
||||
tokio::spawn(async move {
|
||||
let mut err_reader = tokio::io::BufReader::new(stderr).lines();
|
||||
@@ -123,12 +135,11 @@ impl Tunnel for NgrokTunnel {
|
||||
});
|
||||
}
|
||||
|
||||
if let Ok(mut guard) = self.url.write() {
|
||||
*guard = Some(public_url.clone());
|
||||
}
|
||||
|
||||
let mut guard = self.proc.lock().await;
|
||||
*guard = Some(TunnelProcess { child });
|
||||
*guard = Some(TunnelProcess {
|
||||
child,
|
||||
_pipe_drain: Some(drain_handle),
|
||||
});
|
||||
|
||||
Ok(public_url)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user