Files
cls 9dfa06ffee
Docker image / Build (linux/amd64) (push) Has been cancelled
Docker image / Build (linux/arm64) (push) Has been cancelled
Docker image / Merge release multi-arch manifest (push) Has been cancelled
Docker image / Merge debug multi-arch manifest (push) Has been cancelled
Docker image / Build public push gateway (linux/amd64) (push) Has been cancelled
Docker image / Build public push gateway (linux/arm64) (push) Has been cancelled
Docker image / Publish public push gateway image (push) Has been cancelled
Sprig image / Build (linux/amd64) (push) Has been cancelled
Sprig image / Build (linux/arm64) (push) Has been cancelled
Sprig image / Merge multi-arch manifest (push) Has been cancelled
Harbor Buzz Orchestra / Python tests and lint (push) Has been cancelled
CI / Detect Changed Paths (push) Has been cancelled
CI / Rust Lint (push) Has been cancelled
CI / Unit Tests (push) Has been cancelled
CI / Desktop Core (push) Has been cancelled
CI / Desktop Smoke E2E (1) (push) Has been cancelled
CI / Desktop Smoke E2E (2) (push) Has been cancelled
CI / Desktop Smoke E2E (3) (push) Has been cancelled
CI / Desktop Smoke E2E (4) (push) Has been cancelled
CI / Desktop (push) Has been cancelled
CI / Desktop E2E Relay (push) Has been cancelled
CI / Desktop E2E Integration (1/2) (push) Has been cancelled
CI / Desktop E2E Integration (2/2) (push) Has been cancelled
CI / Desktop E2E Integration (push) Has been cancelled
CI / Backend Integration (relay e2e) (push) Has been cancelled
CI / Relay E2E (push) Has been cancelled
CI / Web (push) Has been cancelled
CI / Mobile (push) Has been cancelled
CI / Security (push) Has been cancelled
CI / Dead Token Reference Guard (push) Has been cancelled
CI / Server Cross-Compile (aarch64-unknown-linux-musl) (push) Has been cancelled
CI / Server Cross-Compile (x86_64-unknown-linux-musl) (push) Has been cancelled
CI / Windows Rust (x86_64-pc-windows-msvc) (push) Has been cancelled
CI / Desktop Build (macOS) (push) Has been cancelled
helm chart / lint + unittest + render matrix (push) Has been cancelled
helm chart / install on kind (gated) (push) Has been cancelled
helm chart / publish chart to GHCR (push) Has been cancelled
Mesh Lifecycle / Relay-Driven Mesh Lifecycle Smoke (push) Has been cancelled
Sprig / Build (aarch64-unknown-linux-musl) (push) Has been cancelled
Sprig / Build (x86_64-unknown-linux-musl) (push) Has been cancelled
Sprig / Publish rolling release (push) Has been cancelled
Sprig / Publish tagged release (push) Has been cancelled
feat: import Chinese-localized Buzz source snapshot
Signed-off-by: cls_宁波本机 <908705107@qq.com>
2026-08-13 18:34:25 +08:00

74 lines
2.7 KiB
Rust

//! Loopback TCP forwarder for the benchmark task container.
//!
//! The benchmark relay is host-header tenant-bound: its community row is the
//! authority of its own `RELAY_URL` (e.g. `localhost:3600`), and a request
//! presenting any other `Host` fails closed. Agents inside a Harbor task
//! container can only reach the host-published relay via the Docker host
//! alias (`host.docker.internal`), which would present the wrong `Host`.
//!
//! So the container runtime uploads this forwarder next to the agent stack:
//! agents dial `ws://localhost:<port>` — presenting the exact `Host` the
//! community row expects — and the forwarder bridges the byte stream to the
//! host gateway. Transparent to everything above TCP (WebSocket, the buzz
//! CLI, git-over-HTTP). std-only; compiled with plain `rustc` against the
//! musl target, so it runs on any Linux task image.
//!
//! Usage: `relay-forwarder <listen-addr> <target-addr>`
use std::io::{self, Read, Write};
use std::net::{Shutdown, TcpListener, TcpStream};
use std::thread;
fn main() -> io::Result<()> {
let mut args = std::env::args().skip(1);
let (listen, target) = match (args.next(), args.next()) {
(Some(listen), Some(target)) => (listen, target),
_ => {
eprintln!("usage: relay-forwarder <listen-addr> <target-addr>");
std::process::exit(2);
}
};
let listener = TcpListener::bind(&listen)?;
// Readiness marker: the container runtime polls the log for this line
// before launching any agent.
println!("forwarding {listen} -> {target}");
for client in listener.incoming() {
let Ok(client) = client else { continue };
let target = target.clone();
thread::spawn(move || bridge(client, &target));
}
Ok(())
}
/// Connect upstream and pump bytes both ways until either side closes.
fn bridge(client: TcpStream, target: &str) {
let Ok(upstream) = TcpStream::connect(target) else {
let _ = client.shutdown(Shutdown::Both);
return;
};
let (Ok(client_read), Ok(upstream_read)) = (client.try_clone(), upstream.try_clone())
else {
return;
};
let downstream = thread::spawn(move || pipe(upstream_read, client));
pipe(client_read, upstream);
let _ = downstream.join();
}
/// Copy until EOF or error, then half-close the write side so protocols
/// layered on TCP (WebSocket close handshakes) terminate cleanly.
fn pipe(mut from: TcpStream, mut to: TcpStream) {
let mut buf = [0u8; 16 * 1024];
loop {
match from.read(&mut buf) {
Ok(0) | Err(_) => break,
Ok(n) => {
if to.write_all(&buf[..n]).is_err() {
break;
}
}
}
}
let _ = to.shutdown(Shutdown::Write);
}