Skip to content

Commit c2572ad

Browse files
committed
feat(exec): add bidirectional streaming for interactive TTY sessions
The existing ExecSandbox RPC sends stdin upfront in the request body, making interactive programs (bash, top, vim) unusable with --tty. This adds an ExecSandboxInteractive bidirectional streaming RPC that forwards live keystrokes, terminal resize events, and EOF to the sandbox. - Add ExecSandboxInteractive RPC, ExecSandboxInput and ExecSandboxWindowResize messages to the proto definition - Implement server handler with relay bridging, PTY allocation via russh channel.split(), and timeout support - Implement CLI client with raw mode, spawn_blocking stdin reader, SIGWINCH resize forwarding, and process::exit for clean shutdown - Route to interactive path only when --tty is explicitly passed Signed-off-by: Florent Benoit <fbenoit@redhat.com>
1 parent df5a8b9 commit c2572ad

15 files changed

Lines changed: 787 additions & 50 deletions

crates/openshell-cli/src/run.rs

Lines changed: 137 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2396,6 +2396,11 @@ pub async fn sandbox_exec_grpc(
23962396
let tty = tty_override
23972397
.unwrap_or_else(|| std::io::stdin().is_terminal() && std::io::stdout().is_terminal());
23982398

2399+
if tty_override == Some(true) && std::io::stdin().is_terminal() {
2400+
return sandbox_exec_interactive_grpc(client, &sandbox, command, workdir, timeout_seconds)
2401+
.await;
2402+
}
2403+
23992404
// Make the streaming gRPC call.
24002405
let mut stream = client
24012406
.exec_sandbox(ExecSandboxRequest {
@@ -2406,6 +2411,7 @@ pub async fn sandbox_exec_grpc(
24062411
timeout_seconds,
24072412
stdin: stdin_payload,
24082413
tty,
2414+
..Default::default()
24092415
})
24102416
.await
24112417
.into_diagnostic()?
@@ -2725,6 +2731,137 @@ async fn drain_and_shutdown_local_socket(mut socket: tokio::net::TcpStream) {
27252731
let _ = socket.shutdown().await;
27262732
}
27272733

2734+
struct RawModeGuard;
2735+
2736+
impl Drop for RawModeGuard {
2737+
fn drop(&mut self) {
2738+
let _ = crossterm::terminal::disable_raw_mode();
2739+
}
2740+
}
2741+
2742+
async fn sandbox_exec_interactive_grpc(
2743+
mut client: crate::tls::GrpcClient,
2744+
sandbox: &Sandbox,
2745+
command: &[String],
2746+
workdir: Option<&str>,
2747+
timeout_seconds: u32,
2748+
) -> Result<i32> {
2749+
use futures::SinkExt;
2750+
use openshell_core::proto::{ExecSandboxInput, ExecSandboxWindowResize, exec_sandbox_input};
2751+
2752+
let (cols, rows) = crossterm::terminal::size().unwrap_or((80, 24));
2753+
2754+
let (mut input_tx, input_rx) = futures::channel::mpsc::channel::<ExecSandboxInput>(4096);
2755+
2756+
// Send the start message with exec metadata.
2757+
input_tx
2758+
.send(ExecSandboxInput {
2759+
payload: Some(exec_sandbox_input::Payload::Start(ExecSandboxRequest {
2760+
sandbox_id: sandbox.object_id().to_string(),
2761+
command: command.to_vec(),
2762+
workdir: workdir.unwrap_or_default().to_string(),
2763+
environment: HashMap::new(),
2764+
timeout_seconds,
2765+
stdin: Vec::new(),
2766+
tty: true,
2767+
cols: u32::from(cols),
2768+
rows: u32::from(rows),
2769+
})),
2770+
})
2771+
.await
2772+
.into_diagnostic()?;
2773+
2774+
let mut stream = client
2775+
.exec_sandbox_interactive(input_rx)
2776+
.await
2777+
.into_diagnostic()?
2778+
.into_inner();
2779+
2780+
// Enable raw mode so keystrokes are forwarded immediately.
2781+
crossterm::terminal::enable_raw_mode().into_diagnostic()?;
2782+
let raw_guard = RawModeGuard;
2783+
2784+
// Stdin reader: read raw bytes one at a time so single keystrokes
2785+
// (e.g. 'q' in top) are forwarded immediately.
2786+
let mut stdin_tx = input_tx.clone();
2787+
tokio::task::spawn_blocking(move || {
2788+
let mut stdin = std::io::stdin().lock();
2789+
let mut buf = [0u8; 4096];
2790+
loop {
2791+
let n = match stdin.read(&mut buf) {
2792+
Ok(0) | Err(_) => break,
2793+
Ok(n) => n,
2794+
};
2795+
if stdin_tx
2796+
.try_send(ExecSandboxInput {
2797+
payload: Some(exec_sandbox_input::Payload::Stdin(buf[..n].to_vec())),
2798+
})
2799+
.is_err()
2800+
{
2801+
break;
2802+
}
2803+
}
2804+
});
2805+
2806+
// SIGWINCH handler: forward terminal resize events.
2807+
#[cfg(unix)]
2808+
let mut resize_tx = input_tx.clone();
2809+
#[cfg(unix)]
2810+
let resize_task = tokio::spawn(async move {
2811+
let mut sig = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::window_change())
2812+
.expect("failed to register SIGWINCH handler");
2813+
while sig.recv().await.is_some() {
2814+
if let Ok((c, r)) = crossterm::terminal::size() {
2815+
let msg = ExecSandboxInput {
2816+
payload: Some(exec_sandbox_input::Payload::Resize(
2817+
ExecSandboxWindowResize {
2818+
cols: u32::from(c),
2819+
rows: u32::from(r),
2820+
},
2821+
)),
2822+
};
2823+
if resize_tx.send(msg).await.is_err() {
2824+
break;
2825+
}
2826+
}
2827+
}
2828+
});
2829+
2830+
let mut exit_code = 0i32;
2831+
let stdout = std::io::stdout();
2832+
2833+
while let Some(event) = stream.next().await {
2834+
let event = event.into_diagnostic()?;
2835+
match event.payload {
2836+
Some(exec_sandbox_event::Payload::Stdout(out)) => {
2837+
let mut handle = stdout.lock();
2838+
handle.write_all(&out.data).into_diagnostic()?;
2839+
handle.flush().into_diagnostic()?;
2840+
}
2841+
Some(exec_sandbox_event::Payload::Stderr(err)) => {
2842+
let mut handle = stdout.lock();
2843+
handle.write_all(&err.data).into_diagnostic()?;
2844+
handle.flush().into_diagnostic()?;
2845+
}
2846+
Some(exec_sandbox_event::Payload::Exit(exit)) => {
2847+
exit_code = exit.exit_code;
2848+
break;
2849+
}
2850+
None => {}
2851+
}
2852+
}
2853+
2854+
#[cfg(unix)]
2855+
resize_task.abort();
2856+
2857+
// Drop the raw mode guard to restore the terminal before exiting.
2858+
drop(raw_guard);
2859+
2860+
// The spawn_blocking stdin reader is stuck on stdin.read() and cannot be
2861+
// cancelled. Force-exit so the tokio runtime doesn't hang waiting for it.
2862+
std::process::exit(exit_code)
2863+
}
2864+
27282865
/// Print a single YAML line with dimmed keys and regular values.
27292866
fn print_yaml_line(line: &str) {
27302867
// Find leading whitespace

crates/openshell-cli/tests/ensure_providers_integration.rs

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -13,14 +13,15 @@ use openshell_core::proto::{
1313
CreateSandboxRequest, CreateSshSessionRequest, CreateSshSessionResponse, DeleteProviderRequest,
1414
DeleteProviderResponse, DeleteSandboxRequest, DeleteSandboxResponse,
1515
DetachSandboxProviderRequest, DetachSandboxProviderResponse, ExecSandboxEvent,
16-
ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest, GetGatewayConfigResponse,
17-
GetProviderRequest, GetSandboxConfigRequest, GetSandboxConfigResponse,
18-
GetSandboxProviderEnvironmentRequest, GetSandboxProviderEnvironmentResponse, GetSandboxRequest,
19-
HealthRequest, HealthResponse, ListProvidersRequest, ListProvidersResponse,
20-
ListSandboxProvidersRequest, ListSandboxProvidersResponse, ListSandboxesRequest,
21-
ListSandboxesResponse, Provider, ProviderResponse, RevokeSshSessionRequest,
22-
RevokeSshSessionResponse, SandboxResponse, SandboxStreamEvent, ServiceStatus,
23-
SupervisorMessage, UpdateProviderRequest, WatchSandboxRequest,
16+
ExecSandboxInput, ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest,
17+
GetGatewayConfigResponse, GetProviderRequest, GetSandboxConfigRequest,
18+
GetSandboxConfigResponse, GetSandboxProviderEnvironmentRequest,
19+
GetSandboxProviderEnvironmentResponse, GetSandboxRequest, HealthRequest, HealthResponse,
20+
ListProvidersRequest, ListProvidersResponse, ListSandboxProvidersRequest,
21+
ListSandboxProvidersResponse, ListSandboxesRequest, ListSandboxesResponse, Provider,
22+
ProviderResponse, RevokeSshSessionRequest, RevokeSshSessionResponse, SandboxResponse,
23+
SandboxStreamEvent, ServiceStatus, SupervisorMessage, UpdateProviderRequest,
24+
WatchSandboxRequest,
2425
};
2526
use openshell_core::{ObjectId, ObjectName};
2627
use rcgen::{
@@ -395,6 +396,15 @@ impl OpenShell for TestOpenShell {
395396
)))
396397
}
397398

399+
type ExecSandboxInteractiveStream =
400+
tokio_stream::wrappers::ReceiverStream<Result<ExecSandboxEvent, Status>>;
401+
async fn exec_sandbox_interactive(
402+
&self,
403+
_request: tonic::Request<tonic::Streaming<ExecSandboxInput>>,
404+
) -> Result<Response<Self::ExecSandboxInteractiveStream>, Status> {
405+
Err(Status::unimplemented("not implemented in test"))
406+
}
407+
398408
async fn update_config(
399409
&self,
400410
_request: tonic::Request<openshell_core::proto::UpdateConfigRequest>,

crates/openshell-cli/tests/mtls_integration.rs

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,10 +4,10 @@
44
use openshell_cli::tls::{TlsOptions, grpc_client};
55
use openshell_core::proto::{
66
CreateProviderRequest, CreateSshSessionRequest, CreateSshSessionResponse,
7-
DeleteProviderRequest, DeleteProviderResponse, ExecSandboxEvent, ExecSandboxRequest,
8-
GetProviderRequest, HealthRequest, HealthResponse, ListProvidersRequest, ListProvidersResponse,
9-
ProviderResponse, RevokeSshSessionRequest, RevokeSshSessionResponse, ServiceStatus,
10-
UpdateProviderRequest,
7+
DeleteProviderRequest, DeleteProviderResponse, ExecSandboxEvent, ExecSandboxInput,
8+
ExecSandboxRequest, GetProviderRequest, HealthRequest, HealthResponse, ListProvidersRequest,
9+
ListProvidersResponse, ProviderResponse, RevokeSshSessionRequest, RevokeSshSessionResponse,
10+
ServiceStatus, UpdateProviderRequest,
1111
open_shell_server::{OpenShell, OpenShellServer},
1212
};
1313
use rcgen::{
@@ -286,6 +286,15 @@ impl OpenShell for TestOpenShell {
286286
)))
287287
}
288288

289+
type ExecSandboxInteractiveStream =
290+
tokio_stream::wrappers::ReceiverStream<Result<ExecSandboxEvent, Status>>;
291+
async fn exec_sandbox_interactive(
292+
&self,
293+
_request: tonic::Request<tonic::Streaming<ExecSandboxInput>>,
294+
) -> Result<Response<Self::ExecSandboxInteractiveStream>, Status> {
295+
Err(Status::unimplemented("not implemented in test"))
296+
}
297+
289298
async fn update_config(
290299
&self,
291300
_request: tonic::Request<openshell_core::proto::UpdateConfigRequest>,

crates/openshell-cli/tests/provider_commands_integration.rs

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -9,14 +9,15 @@ use openshell_core::proto::{
99
CreateSandboxRequest, CreateSshSessionRequest, CreateSshSessionResponse, DeleteProviderRequest,
1010
DeleteProviderResponse, DeleteSandboxRequest, DeleteSandboxResponse,
1111
DetachSandboxProviderRequest, DetachSandboxProviderResponse, ExecSandboxEvent,
12-
ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest, GetGatewayConfigResponse,
13-
GetProviderRequest, GetSandboxConfigRequest, GetSandboxConfigResponse,
14-
GetSandboxProviderEnvironmentRequest, GetSandboxProviderEnvironmentResponse, GetSandboxRequest,
15-
HealthRequest, HealthResponse, ListProvidersRequest, ListProvidersResponse,
16-
ListSandboxProvidersRequest, ListSandboxProvidersResponse, ListSandboxesRequest,
17-
ListSandboxesResponse, Provider, ProviderProfile, ProviderResponse, RevokeSshSessionRequest,
18-
RevokeSshSessionResponse, SandboxResponse, SandboxStreamEvent, ServiceStatus,
19-
SupervisorMessage, UpdateProviderRequest, WatchSandboxRequest,
12+
ExecSandboxInput, ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest,
13+
GetGatewayConfigResponse, GetProviderRequest, GetSandboxConfigRequest,
14+
GetSandboxConfigResponse, GetSandboxProviderEnvironmentRequest,
15+
GetSandboxProviderEnvironmentResponse, GetSandboxRequest, HealthRequest, HealthResponse,
16+
ListProvidersRequest, ListProvidersResponse, ListSandboxProvidersRequest,
17+
ListSandboxProvidersResponse, ListSandboxesRequest, ListSandboxesResponse, Provider,
18+
ProviderProfile, ProviderResponse, RevokeSshSessionRequest, RevokeSshSessionResponse,
19+
SandboxResponse, SandboxStreamEvent, ServiceStatus, SupervisorMessage, UpdateProviderRequest,
20+
WatchSandboxRequest,
2021
};
2122
use openshell_core::{ObjectId, ObjectName};
2223
use rcgen::{
@@ -504,6 +505,15 @@ impl OpenShell for TestOpenShell {
504505
)))
505506
}
506507

508+
type ExecSandboxInteractiveStream =
509+
tokio_stream::wrappers::ReceiverStream<Result<ExecSandboxEvent, Status>>;
510+
async fn exec_sandbox_interactive(
511+
&self,
512+
_request: tonic::Request<tonic::Streaming<ExecSandboxInput>>,
513+
) -> Result<Response<Self::ExecSandboxInteractiveStream>, Status> {
514+
Err(Status::unimplemented("not implemented in test"))
515+
}
516+
507517
async fn update_config(
508518
&self,
509519
_request: tonic::Request<openshell_core::proto::UpdateConfigRequest>,

crates/openshell-cli/tests/sandbox_create_lifecycle_integration.rs

Lines changed: 18 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -10,15 +10,15 @@ use openshell_core::proto::{
1010
CreateSandboxRequest, CreateSshSessionRequest, CreateSshSessionResponse, DeleteProviderRequest,
1111
DeleteProviderResponse, DeleteSandboxRequest, DeleteSandboxResponse,
1212
DetachSandboxProviderRequest, DetachSandboxProviderResponse, ExecSandboxEvent,
13-
ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest, GetGatewayConfigResponse,
14-
GetProviderRequest, GetSandboxConfigRequest, GetSandboxConfigResponse,
15-
GetSandboxProviderEnvironmentRequest, GetSandboxProviderEnvironmentResponse, GetSandboxRequest,
16-
HealthRequest, HealthResponse, ListProvidersRequest, ListProvidersResponse,
17-
ListSandboxProvidersRequest, ListSandboxProvidersResponse, ListSandboxesRequest,
18-
ListSandboxesResponse, PlatformEvent, ProviderResponse, RevokeSshSessionRequest,
19-
RevokeSshSessionResponse, Sandbox, SandboxPhase, SandboxResponse, SandboxStreamEvent,
20-
ServiceStatus, SupervisorMessage, UpdateProviderRequest, WatchSandboxRequest,
21-
sandbox_stream_event,
13+
ExecSandboxInput, ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest,
14+
GetGatewayConfigResponse, GetProviderRequest, GetSandboxConfigRequest,
15+
GetSandboxConfigResponse, GetSandboxProviderEnvironmentRequest,
16+
GetSandboxProviderEnvironmentResponse, GetSandboxRequest, HealthRequest, HealthResponse,
17+
ListProvidersRequest, ListProvidersResponse, ListSandboxProvidersRequest,
18+
ListSandboxProvidersResponse, ListSandboxesRequest, ListSandboxesResponse, PlatformEvent,
19+
ProviderResponse, RevokeSshSessionRequest, RevokeSshSessionResponse, Sandbox, SandboxPhase,
20+
SandboxResponse, SandboxStreamEvent, ServiceStatus, SupervisorMessage, UpdateProviderRequest,
21+
WatchSandboxRequest, sandbox_stream_event,
2222
};
2323
use rcgen::{
2424
BasicConstraints, Certificate, CertificateParams, ExtendedKeyUsagePurpose, IsCa, KeyPair,
@@ -369,6 +369,15 @@ impl OpenShell for TestOpenShell {
369369
)))
370370
}
371371

372+
type ExecSandboxInteractiveStream =
373+
tokio_stream::wrappers::ReceiverStream<Result<ExecSandboxEvent, Status>>;
374+
async fn exec_sandbox_interactive(
375+
&self,
376+
_request: tonic::Request<tonic::Streaming<ExecSandboxInput>>,
377+
) -> Result<Response<Self::ExecSandboxInteractiveStream>, Status> {
378+
Err(Status::unimplemented("not implemented in test"))
379+
}
380+
372381
async fn update_config(
373382
&self,
374383
_request: tonic::Request<openshell_core::proto::UpdateConfigRequest>,

crates/openshell-cli/tests/sandbox_name_fallback_integration.rs

Lines changed: 17 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -10,14 +10,14 @@ use openshell_core::proto::{
1010
CreateSandboxRequest, CreateSshSessionRequest, CreateSshSessionResponse, DeleteProviderRequest,
1111
DeleteProviderResponse, DeleteSandboxRequest, DeleteSandboxResponse,
1212
DetachSandboxProviderRequest, DetachSandboxProviderResponse, ExecSandboxEvent,
13-
ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest, GetGatewayConfigResponse,
14-
GetProviderRequest, GetSandboxConfigRequest, GetSandboxConfigResponse,
15-
GetSandboxProviderEnvironmentRequest, GetSandboxProviderEnvironmentResponse, GetSandboxRequest,
16-
HealthRequest, HealthResponse, ListProvidersRequest, ListProvidersResponse,
17-
ListSandboxProvidersRequest, ListSandboxProvidersResponse, ListSandboxesRequest,
18-
ListSandboxesResponse, ProviderResponse, Sandbox, SandboxPolicy, SandboxResponse,
19-
SandboxStreamEvent, ServiceStatus, SupervisorMessage, UpdateProviderRequest,
20-
WatchSandboxRequest,
13+
ExecSandboxInput, ExecSandboxRequest, GatewayMessage, GetGatewayConfigRequest,
14+
GetGatewayConfigResponse, GetProviderRequest, GetSandboxConfigRequest,
15+
GetSandboxConfigResponse, GetSandboxProviderEnvironmentRequest,
16+
GetSandboxProviderEnvironmentResponse, GetSandboxRequest, HealthRequest, HealthResponse,
17+
ListProvidersRequest, ListProvidersResponse, ListSandboxProvidersRequest,
18+
ListSandboxProvidersResponse, ListSandboxesRequest, ListSandboxesResponse, ProviderResponse,
19+
Sandbox, SandboxPolicy, SandboxResponse, SandboxStreamEvent, ServiceStatus, SupervisorMessage,
20+
UpdateProviderRequest, WatchSandboxRequest,
2121
};
2222
use rcgen::{
2323
BasicConstraints, Certificate, CertificateParams, ExtendedKeyUsagePurpose, IsCa, KeyPair,
@@ -307,6 +307,15 @@ impl OpenShell for TestOpenShell {
307307
)))
308308
}
309309

310+
type ExecSandboxInteractiveStream =
311+
tokio_stream::wrappers::ReceiverStream<Result<ExecSandboxEvent, Status>>;
312+
async fn exec_sandbox_interactive(
313+
&self,
314+
_request: tonic::Request<tonic::Streaming<ExecSandboxInput>>,
315+
) -> Result<Response<Self::ExecSandboxInteractiveStream>, Status> {
316+
Err(Status::unimplemented("not implemented in test"))
317+
}
318+
310319
async fn update_config(
311320
&self,
312321
_request: tonic::Request<openshell_core::proto::UpdateConfigRequest>,

crates/openshell-server/src/grpc/mod.rs

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,10 +16,10 @@ use openshell_core::proto::{
1616
DeleteProviderProfileResponse, DeleteProviderRequest, DeleteProviderResponse,
1717
DeleteSandboxRequest, DeleteSandboxResponse, DetachSandboxProviderRequest,
1818
DetachSandboxProviderResponse, EditDraftChunkRequest, EditDraftChunkResponse, ExecSandboxEvent,
19-
ExecSandboxRequest, GatewayMessage, GetDraftHistoryRequest, GetDraftHistoryResponse,
20-
GetDraftPolicyRequest, GetDraftPolicyResponse, GetGatewayConfigRequest,
21-
GetGatewayConfigResponse, GetProviderProfileRequest, GetProviderRequest,
22-
GetSandboxConfigRequest, GetSandboxConfigResponse, GetSandboxLogsRequest,
19+
ExecSandboxInput, ExecSandboxRequest, GatewayMessage, GetDraftHistoryRequest,
20+
GetDraftHistoryResponse, GetDraftPolicyRequest, GetDraftPolicyResponse,
21+
GetGatewayConfigRequest, GetGatewayConfigResponse, GetProviderProfileRequest,
22+
GetProviderRequest, GetSandboxConfigRequest, GetSandboxConfigResponse, GetSandboxLogsRequest,
2323
GetSandboxLogsResponse, GetSandboxPolicyStatusRequest, GetSandboxPolicyStatusResponse,
2424
GetSandboxProviderEnvironmentRequest, GetSandboxProviderEnvironmentResponse, GetSandboxRequest,
2525
HealthRequest, HealthResponse, ImportProviderProfilesRequest, ImportProviderProfilesResponse,
@@ -251,6 +251,15 @@ impl OpenShell for OpenShellService {
251251
sandbox::handle_forward_tcp(&self.state, request).await
252252
}
253253

254+
type ExecSandboxInteractiveStream = ReceiverStream<Result<ExecSandboxEvent, Status>>;
255+
256+
async fn exec_sandbox_interactive(
257+
&self,
258+
request: Request<tonic::Streaming<ExecSandboxInput>>,
259+
) -> Result<Response<Self::ExecSandboxInteractiveStream>, Status> {
260+
sandbox::handle_exec_sandbox_interactive(&self.state, request).await
261+
}
262+
254263
// --- SSH sessions ---
255264

256265
async fn create_ssh_session(

0 commit comments

Comments
 (0)