Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 27 additions & 1 deletion architecture/gateway.md
Original file line number Diff line number Diff line change
Expand Up @@ -1137,7 +1137,33 @@ name, an authorization token from `CreateSshSession`, and an explicit target:
`target.ssh` for the sandbox SSH socket or `target.tcp` for a loopback service
inside the sandbox. The gateway validates the token and sandbox readiness,
sends a targeted `RelayOpen` to the supervisor, then bridges
`TcpForwardFrame::Data` to `RelayFrame::Data` until either side closes.
`TcpForwardFrame::Data` to `RelayFrame::Data`. Request EOF closes only that
input direction. Peers advertising `stream-half-close-v1` receive a response
FIN after response bytes and keep the request direction open until its own EOF.
Final success follows both directions; consumers must still check trailers.
Legacy peers receive response EOF and retain the old completion behavior.

Each forwarding hop negotiates independently, including gateway-to-gateway
`PeerRelay`, and propagates downstream limitations upstream. PeerRelay response
metadata confirms effective support; absence selects legacy EOF for the path.
Supervisor capabilities are exchanged through Hello/Accepted and
carried with the owning session ID in RelayInit. The SSH stdio proxy uses legacy
response EOF because stdout cannot represent an independent socket FIN.

Bridge tasks own both directional pumps. Internal relay pipes retain downstream
RPC completion separately from byte EOF, so forwarding FIN does not hide later
error trailers. They retain typed abort status and wake blocked readers,
writers, and frame senders, so failure
cannot become a successful byte EOF. `RelayClose` is an abort, scoped to its
sandbox and supervisor session; stale or unrelated sessions cannot cancel an
active channel. A close from the current session also removes an owned pending
channel and delivers its typed error to the waiting caller before claim, freeing
pending capacity immediately. The first observed abort wins, with transport
failure as fallback when a typed reason cannot be delivered. Control-session loss
alone does not cancel established data relays. Those retain their original session ownership;
replacement sessions cannot cancel them. Legacy reverse relays close their
response stream on input EOF while still draining delayed target replies.
No new application-idle or half-closed timeout is imposed.

Browser service URLs use the same supervisor relay path after host-based
routing resolves `sandbox--service.<service-routing-domain>` to a stored
Expand Down
51 changes: 5 additions & 46 deletions crates/openshell-cli/src/run.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2184,7 +2184,6 @@ async fn forward_one_tcp_connection(
service_id: String,
authorization_token: String,
) -> std::result::Result<(), ForwardTcpConnectionError> {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio_stream::wrappers::ReceiverStream;

let (tx, rx) = tokio::sync::mpsc::channel::<TcpForwardFrame>(16);
Expand All @@ -2199,13 +2198,14 @@ async fn forward_one_tcp_connection(
port: u32::from(target_port),
})),
authorization_token,
capabilities: openshell_core::stream_lifecycle::capabilities(),
},
)),
})
.await
.map_err(|_| ForwardTcpConnectionError::transient("failed to initialize forward stream"))?;

let mut response = match client.forward_tcp(ReceiverStream::new(rx)).await {
let response = match client.forward_tcp(ReceiverStream::new(rx)).await {
Ok(response) => response.into_inner(),
Err(status) => {
let err = ForwardTcpConnectionError::from_status(status);
Expand All @@ -2214,51 +2214,10 @@ async fn forward_one_tcp_connection(
}
};

let (mut local_read, mut local_write) = socket.into_split();

let to_gateway = tokio::spawn(async move {
let mut buf = vec![0u8; 64 * 1024];
loop {
let n = local_read.read(&mut buf).await?;
if n == 0 {
break;
}
if tx
.send(TcpForwardFrame {
payload: Some(openshell_core::proto::tcp_forward_frame::Payload::Data(
buf[..n].to_vec(),
)),
})
.await
.is_err()
{
break;
}
}
Ok::<(), std::io::Error>(())
});

while let Some(frame) = response
.message()
let (local_read, local_write) = socket.into_split();
openshell_core::stream_lifecycle::client(response, local_read, local_write, tx, true)
.await
.map_err(ForwardTcpConnectionError::from_status)?
{
let Some(openshell_core::proto::tcp_forward_frame::Payload::Data(data)) = frame.payload
else {
continue;
};
if data.is_empty() {
continue;
}
local_write
.write_all(&data)
.await
.map_err(|err| ForwardTcpConnectionError::transient(err.to_string()))?;
}

let _ = local_write.shutdown().await;
to_gateway.abort();
Ok(())
.map_err(ForwardTcpConnectionError::from_status)
}

async fn drain_and_shutdown_local_socket(mut socket: tokio::net::TcpStream) {
Expand Down
59 changes: 11 additions & 48 deletions crates/openshell-cli/src/ssh.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@ use std::os::unix::process::CommandExt;
use std::path::{Path, PathBuf};
use std::process::{Command, ExitStatus, Stdio};
use std::time::{Duration, Instant};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpStream;
use tokio::process::{Child, Command as TokioCommand};
use tokio_stream::wrappers::ReceiverStream;
Expand Down Expand Up @@ -1914,64 +1913,28 @@ pub async fn sandbox_ssh_proxy(
service_id: format!("ssh-proxy:{sandbox_name}"),
target: Some(tcp_forward_init::Target::Ssh(SshRelayTarget {})),
authorization_token: token.to_string(),
capabilities: Vec::new(),
},
)),
})
.await
.map_err(|_| miette::miette!("failed to initialize SSH forward stream"))?;

let mut response = client
let response = client
.forward_tcp(ReceiverStream::new(rx))
.await
.into_diagnostic()?
.into_inner();

let stdin = tokio::io::stdin();
let stdout = tokio::io::stdout();

let to_remote = tokio::spawn(async move {
let mut stdin = stdin;
let mut buf = vec![0u8; 64 * 1024];
while let Ok(n) = stdin.read(&mut buf).await {
if n == 0 {
break;
}
if tx
.send(TcpForwardFrame {
payload: Some(openshell_core::proto::tcp_forward_frame::Payload::Data(
buf[..n].to_vec(),
)),
})
.await
.is_err()
{
break;
}
}
});
let from_remote = tokio::spawn(async move {
let mut stdout = stdout;
loop {
let Ok(Some(frame)) = response.message().await else {
break;
};
let Some(openshell_core::proto::tcp_forward_frame::Payload::Data(data)) = frame.payload
else {
continue;
};
if data.is_empty() {
continue;
}
if stdout.write_all(&data).await.is_err() {
break;
}
let _ = stdout.flush().await;
}
});
let _ = from_remote.await;
to_remote.abort();

Ok(())
openshell_core::stream_lifecycle::client(
response,
tokio::io::stdin(),
tokio::io::stdout(),
tx,
false,
)
.await
.into_diagnostic()
}

fn grpc_server_from_ssh_gateway_url(gateway_url: &str) -> Result<String> {
Expand Down
1 change: 1 addition & 0 deletions crates/openshell-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ pub mod secrets;
pub mod settings;
pub mod shell;
pub mod spiffe;
pub mod stream_lifecycle;
pub mod telemetry;
pub mod time;
pub mod transport_errors;
Expand Down
Loading
Loading