@@ -1852,6 +1852,86 @@ fn local_terminal_size() -> Option<(u32, u32)> {
18521852const MAX_EXEC_REQUEST_BYTES : usize = 1024 * 1024 ;
18531853const MAX_EXEC_STDIN_BYTES : usize = 4 * 1024 * 1024 ;
18541854
1855+ /// How long `sandbox exec` waits for piped stdin to reach EOF before it starts
1856+ /// the command and streams the rest of the input as it arrives.
1857+ ///
1858+ /// Small pipes such as `echo x | openshell sandbox exec ...` close well within
1859+ /// this window and keep using the single-request path that older gateways
1860+ /// need. A pipe that stays open (a CI runner, a supervisor, a harness that
1861+ /// never closes stdin) must not block the command: after the grace period the
1862+ /// command starts and stdin is forwarded until EOF or until the command exits.
1863+ const EXEC_STDIN_UNARY_GRACE : Duration = Duration :: from_millis ( 200 ) ;
1864+
1865+ /// Piped stdin as collected by [`collect_piped_stdin`].
1866+ enum PipedStdin {
1867+ /// stdin reached EOF within the grace period; `prefix` holds all of it.
1868+ Complete ( Vec < u8 > ) ,
1869+ /// stdin is still open. `prefix` is what arrived so far; `rest` delivers
1870+ /// the remaining chunks until EOF.
1871+ Open {
1872+ prefix : Vec < u8 > ,
1873+ rest : tokio:: sync:: mpsc:: Receiver < std:: io:: Result < Vec < u8 > > > ,
1874+ } ,
1875+ }
1876+
1877+ /// Read a pipe on a detached OS thread and hand its bytes over in chunks.
1878+ ///
1879+ /// A plain `std::thread` rather than `spawn_blocking`, so runtime shutdown
1880+ /// never waits on a thread parked in `read(2)`. The thread exits at EOF, on a
1881+ /// read error, or when the receiver is dropped.
1882+ fn spawn_piped_stdin_reader (
1883+ mut reader : impl Read + Send + ' static ,
1884+ ) -> tokio:: sync:: mpsc:: Receiver < std:: io:: Result < Vec < u8 > > > {
1885+ let ( tx, rx) = tokio:: sync:: mpsc:: channel :: < std:: io:: Result < Vec < u8 > > > ( 64 ) ;
1886+ std:: thread:: spawn ( move || {
1887+ let mut buf = [ 0u8 ; 4096 ] ;
1888+ loop {
1889+ match reader. read ( & mut buf) {
1890+ Ok ( 0 ) => return ,
1891+ Err ( error) if error. kind ( ) == ErrorKind :: Interrupted => { }
1892+ Err ( error) => {
1893+ let _ = tx. blocking_send ( Err ( error) ) ;
1894+ return ;
1895+ }
1896+ Ok ( n) => {
1897+ if tx. blocking_send ( Ok ( buf[ ..n] . to_vec ( ) ) ) . is_err ( ) {
1898+ return ;
1899+ }
1900+ }
1901+ }
1902+ }
1903+ } ) ;
1904+ rx
1905+ }
1906+
1907+ /// Collect piped stdin until EOF or until `grace` elapses, whichever comes
1908+ /// first. Input beyond `limit` bytes is rejected with the upload hint.
1909+ async fn collect_piped_stdin (
1910+ mut rx : tokio:: sync:: mpsc:: Receiver < std:: io:: Result < Vec < u8 > > > ,
1911+ grace : Duration ,
1912+ limit : usize ,
1913+ ) -> Result < PipedStdin > {
1914+ let deadline = tokio:: time:: Instant :: now ( ) + grace;
1915+ let mut prefix = Vec :: new ( ) ;
1916+ loop {
1917+ match tokio:: time:: timeout_at ( deadline, rx. recv ( ) ) . await {
1918+ Ok ( Some ( Ok ( chunk) ) ) => {
1919+ prefix. extend_from_slice ( & chunk) ;
1920+ if prefix. len ( ) > limit {
1921+ return Err ( piped_stdin_limit_error ( ) ) ;
1922+ }
1923+ }
1924+ Ok ( Some ( Err ( error) ) ) => return Err ( error) . into_diagnostic ( ) ,
1925+ Ok ( None ) => return Ok ( PipedStdin :: Complete ( prefix) ) ,
1926+ Err ( _elapsed) => return Ok ( PipedStdin :: Open { prefix, rest : rx } ) ,
1927+ }
1928+ }
1929+ }
1930+
1931+ fn piped_stdin_limit_error ( ) -> miette:: Report {
1932+ miette:: miette!( "piped stdin exceeds the 4 MiB limit; use `sandbox upload` for larger input" )
1933+ }
1934+
18551935/// Execute a command in a running sandbox via gRPC, streaming output to the terminal.
18561936///
18571937/// Returns the remote command's exit code, or an error if the event stream
@@ -1902,24 +1982,24 @@ pub async fn sandbox_exec_grpc(
19021982 // interactive RPC closes the SSH channel when stdin reaches EOF. Retain
19031983 // the existing 4 MiB input cap because the supervisor's process stdin
19041984 // queue is unbounded; larger input should use file upload instead.
1905- let stdin_prefix = if stdin_is_terminal {
1906- Vec :: new ( )
1985+ //
1986+ // Never block on stdin EOF before starting the command: a pipe that stays
1987+ // open (CI runners, harnesses) would otherwise hang the exec forever
1988+ // without the gateway ever seeing the request. After a short grace period
1989+ // the command starts and the remaining input streams until EOF.
1990+ let ( stdin_prefix, stdin_rest) = if stdin_is_terminal {
1991+ ( Vec :: new ( ) , None )
19071992 } else {
1908- tokio:: task:: spawn_blocking ( || {
1909- let mut prefix = Vec :: new ( ) ;
1910- std:: io:: stdin ( )
1911- . take ( ( MAX_EXEC_STDIN_BYTES + 1 ) as u64 )
1912- . read_to_end ( & mut prefix)
1913- . into_diagnostic ( ) ?;
1914- if prefix. len ( ) > MAX_EXEC_STDIN_BYTES {
1915- return Err ( miette:: miette!(
1916- "piped stdin exceeds the 4 MiB limit; use `sandbox upload` for larger input"
1917- ) ) ;
1918- }
1919- Ok :: < _ , miette:: Report > ( prefix)
1920- } )
1921- . await
1922- . into_diagnostic ( ) ??
1993+ match collect_piped_stdin (
1994+ spawn_piped_stdin_reader ( std:: io:: stdin ( ) ) ,
1995+ EXEC_STDIN_UNARY_GRACE ,
1996+ MAX_EXEC_STDIN_BYTES ,
1997+ )
1998+ . await ?
1999+ {
2000+ PipedStdin :: Complete ( prefix) => ( prefix, None ) ,
2001+ PipedStdin :: Open { prefix, rest } => ( prefix, Some ( rest) ) ,
2002+ }
19232003 } ;
19242004
19252005 let ( cols, rows) = if tty {
@@ -1952,7 +2032,10 @@ pub async fn sandbox_exec_grpc(
19522032 "exec command or environment exceeds the gateway's 1 MiB message limit"
19532033 ) ) ;
19542034 }
1955- if ( tty && stdin_is_terminal) || request. encoded_len ( ) > MAX_EXEC_REQUEST_BYTES {
2035+ if ( tty && stdin_is_terminal)
2036+ || stdin_rest. is_some ( )
2037+ || request. encoded_len ( ) > MAX_EXEC_REQUEST_BYTES
2038+ {
19562039 return sandbox_exec_streaming_grpc (
19572040 client,
19582041 & sandbox,
@@ -1964,6 +2047,7 @@ pub async fn sandbox_exec_grpc(
19642047 tty,
19652048 stdin_is_terminal,
19662049 std:: mem:: take ( & mut request. stdin ) ,
2050+ stdin_rest,
19672051 )
19682052 . await ;
19692053 }
@@ -2350,6 +2434,7 @@ async fn sandbox_exec_streaming_grpc(
23502434 tty : bool ,
23512435 stdin_is_terminal : bool ,
23522436 stdin_prefix : Vec < u8 > ,
2437+ stdin_rest : Option < tokio:: sync:: mpsc:: Receiver < std:: io:: Result < Vec < u8 > > > > ,
23532438) -> Result < i32 > {
23542439 #[ cfg( unix) ]
23552440 use openshell_core:: proto:: ExecSandboxWindowResize ;
@@ -2407,42 +2492,88 @@ async fn sandbox_exec_streaming_grpc(
24072492 // closes (blocking_send returns Err) or stdin hits EOF.
24082493 let stdin_tx = input_tx. clone ( ) ;
24092494 let ( stdin_result_tx, mut stdin_result_rx) = tokio:: sync:: oneshot:: channel ( ) ;
2410- std:: thread:: spawn ( move || {
2411- let mut stdin = std:: io:: stdin ( ) . lock ( ) ;
2412- let mut buf = [ 0u8 ; 4096 ] ;
2413- let result = ( || {
2414- for chunk in stdin_prefix. chunks ( buf. len ( ) ) {
2415- if stdin_tx
2416- . blocking_send ( ExecSandboxInput {
2417- payload : Some ( exec_sandbox_input:: Payload :: Stdin ( chunk. to_vec ( ) ) ) ,
2418- } )
2419- . is_err ( )
2420- {
2421- return Ok ( ( ) ) ;
2495+ if let Some ( mut rest) = stdin_rest {
2496+ // Piped stdin that was still open when the command started: the
2497+ // reader thread from `spawn_piped_stdin_reader` already owns stdin,
2498+ // so forward its chunks here after the collected prefix. The 4 MiB
2499+ // cap covers the prefix and the streamed remainder together.
2500+ tokio:: spawn ( async move {
2501+ let mut forwarded = 0usize ;
2502+ let result = async {
2503+ for chunk in stdin_prefix. chunks ( 4096 ) {
2504+ forwarded += chunk. len ( ) ;
2505+ if stdin_tx
2506+ . send ( ExecSandboxInput {
2507+ payload : Some ( exec_sandbox_input:: Payload :: Stdin ( chunk. to_vec ( ) ) ) ,
2508+ } )
2509+ . await
2510+ . is_err ( )
2511+ {
2512+ return Ok ( ( ) ) ;
2513+ }
24222514 }
2515+ while let Some ( chunk) = rest. recv ( ) . await {
2516+ let chunk = chunk?;
2517+ forwarded += chunk. len ( ) ;
2518+ if forwarded > MAX_EXEC_STDIN_BYTES {
2519+ return Err ( std:: io:: Error :: new (
2520+ ErrorKind :: InvalidInput ,
2521+ piped_stdin_limit_error ( ) . to_string ( ) ,
2522+ ) ) ;
2523+ }
2524+ if stdin_tx
2525+ . send ( ExecSandboxInput {
2526+ payload : Some ( exec_sandbox_input:: Payload :: Stdin ( chunk) ) ,
2527+ } )
2528+ . await
2529+ . is_err ( )
2530+ {
2531+ return Ok ( ( ) ) ;
2532+ }
2533+ }
2534+ Ok ( ( ) )
24232535 }
2424- loop {
2425- match stdin. read ( & mut buf) {
2426- Ok ( 0 ) => return Ok ( ( ) ) ,
2427- Err ( error) if error. kind ( ) == ErrorKind :: Interrupted => { }
2428- Err ( error) => return Err ( error) ,
2429- Ok ( n) => {
2430- if stdin_tx
2431- . blocking_send ( ExecSandboxInput {
2432- payload : Some ( exec_sandbox_input:: Payload :: Stdin (
2433- buf[ ..n] . to_vec ( ) ,
2434- ) ) ,
2435- } )
2436- . is_err ( )
2437- {
2438- return Ok ( ( ) ) ;
2536+ . await ;
2537+ let _ = stdin_result_tx. send ( result) ;
2538+ } ) ;
2539+ } else {
2540+ std:: thread:: spawn ( move || {
2541+ let mut stdin = std:: io:: stdin ( ) . lock ( ) ;
2542+ let mut buf = [ 0u8 ; 4096 ] ;
2543+ let result = ( || {
2544+ for chunk in stdin_prefix. chunks ( buf. len ( ) ) {
2545+ if stdin_tx
2546+ . blocking_send ( ExecSandboxInput {
2547+ payload : Some ( exec_sandbox_input:: Payload :: Stdin ( chunk. to_vec ( ) ) ) ,
2548+ } )
2549+ . is_err ( )
2550+ {
2551+ return Ok ( ( ) ) ;
2552+ }
2553+ }
2554+ loop {
2555+ match stdin. read ( & mut buf) {
2556+ Ok ( 0 ) => return Ok ( ( ) ) ,
2557+ Err ( error) if error. kind ( ) == ErrorKind :: Interrupted => { }
2558+ Err ( error) => return Err ( error) ,
2559+ Ok ( n) => {
2560+ if stdin_tx
2561+ . blocking_send ( ExecSandboxInput {
2562+ payload : Some ( exec_sandbox_input:: Payload :: Stdin (
2563+ buf[ ..n] . to_vec ( ) ,
2564+ ) ) ,
2565+ } )
2566+ . is_err ( )
2567+ {
2568+ return Ok ( ( ) ) ;
2569+ }
24392570 }
24402571 }
24412572 }
2442- }
2443- } ) ( ) ;
2444- let _ = stdin_result_tx . send ( result ) ;
2445- } ) ;
2573+ } ) ( ) ;
2574+ let _ = stdin_result_tx . send ( result ) ;
2575+ } ) ;
2576+ }
24462577
24472578 // SIGWINCH handler: forward terminal resize events.
24482579 #[ cfg( unix) ]
@@ -8206,4 +8337,82 @@ mod tests {
82068337 let log = log_line ( "OCSF" , "ocsf" , message, "sandbox" , & [ ] ) ;
82078338 assert ! ( format_log_line( & log) . ends_with( message) ) ;
82088339 }
8340+
8341+ use std:: io:: Write as _;
8342+ use std:: time:: { Duration , Instant } ;
8343+
8344+ fn exec_stdin_runtime ( ) -> tokio:: runtime:: Runtime {
8345+ tokio:: runtime:: Builder :: new_current_thread ( )
8346+ . enable_all ( )
8347+ . build ( )
8348+ . expect ( "runtime" )
8349+ }
8350+
8351+ #[ test]
8352+ fn piped_stdin_that_closes_quickly_goes_in_one_request ( ) {
8353+ let ( reader, mut writer) = std:: io:: pipe ( ) . expect ( "pipe" ) ;
8354+ writer. write_all ( b"hi\n " ) . unwrap ( ) ;
8355+ drop ( writer) ;
8356+ let collected = exec_stdin_runtime ( ) . block_on ( super :: collect_piped_stdin (
8357+ super :: spawn_piped_stdin_reader ( reader) ,
8358+ Duration :: from_secs ( 5 ) ,
8359+ super :: MAX_EXEC_STDIN_BYTES ,
8360+ ) ) ;
8361+ match collected. expect ( "collect" ) {
8362+ super :: PipedStdin :: Complete ( prefix) => assert_eq ! ( prefix, b"hi\n " ) ,
8363+ super :: PipedStdin :: Open { .. } => panic ! ( "closed pipe must complete" ) ,
8364+ }
8365+ }
8366+
8367+ #[ test]
8368+ fn piped_stdin_that_stays_open_does_not_block_the_command ( ) {
8369+ let ( reader, mut writer) = std:: io:: pipe ( ) . expect ( "pipe" ) ;
8370+ writer. write_all ( b"early" ) . unwrap ( ) ;
8371+ let grace = Duration :: from_millis ( 100 ) ;
8372+ let started = Instant :: now ( ) ;
8373+ let runtime = exec_stdin_runtime ( ) ;
8374+ let collected = runtime. block_on ( super :: collect_piped_stdin (
8375+ super :: spawn_piped_stdin_reader ( reader) ,
8376+ grace,
8377+ super :: MAX_EXEC_STDIN_BYTES ,
8378+ ) ) ;
8379+ let elapsed = started. elapsed ( ) ;
8380+ assert ! ( elapsed >= grace, "must wait the grace period: {elapsed:?}" ) ;
8381+ assert ! (
8382+ elapsed < Duration :: from_secs( 3 ) ,
8383+ "must not wait for EOF: {elapsed:?}"
8384+ ) ;
8385+ let super :: PipedStdin :: Open { prefix, mut rest } = collected. expect ( "collect" ) else {
8386+ panic ! ( "an open pipe must start the command before EOF" ) ;
8387+ } ;
8388+ assert_eq ! ( prefix, b"early" ) ;
8389+ // The remainder keeps flowing after the command has started.
8390+ writer. write_all ( b"late" ) . unwrap ( ) ;
8391+ drop ( writer) ;
8392+ let next = runtime
8393+ . block_on ( rest. recv ( ) )
8394+ . expect ( "late chunk" )
8395+ . expect ( "read" ) ;
8396+ assert_eq ! ( next, b"late" ) ;
8397+ assert ! (
8398+ runtime. block_on( rest. recv( ) ) . is_none( ) ,
8399+ "EOF closes the channel"
8400+ ) ;
8401+ }
8402+
8403+ #[ test]
8404+ fn piped_stdin_over_the_limit_is_rejected_with_the_upload_hint ( ) {
8405+ let ( reader, mut writer) = std:: io:: pipe ( ) . expect ( "pipe" ) ;
8406+ writer. write_all ( & [ 0u8 ; 8 ] ) . unwrap ( ) ;
8407+ drop ( writer) ;
8408+ let error = exec_stdin_runtime ( )
8409+ . block_on ( super :: collect_piped_stdin (
8410+ super :: spawn_piped_stdin_reader ( reader) ,
8411+ Duration :: from_secs ( 5 ) ,
8412+ 4 ,
8413+ ) )
8414+ . err ( )
8415+ . expect ( "over the limit must fail" ) ;
8416+ assert ! ( error. to_string( ) . contains( "sandbox upload" ) , "{error}" ) ;
8417+ }
82098418}
0 commit comments