44#![ cfg( feature = "e2e" ) ]
55
66use std:: io:: Write ;
7- use std:: process:: Command ;
87use std:: process:: Stdio ;
98use std:: sync:: Mutex ;
10- use std:: time:: Duration ;
119
1210use openshell_e2e:: harness:: binary:: openshell_cmd;
13- use openshell_e2e:: harness:: port:: find_free_port;
1411use openshell_e2e:: harness:: sandbox:: SandboxGuard ;
1512use tempfile:: NamedTempFile ;
16- use tokio:: time:: { interval, timeout} ;
13+ use tokio:: io:: AsyncReadExt ;
14+ use tokio:: io:: AsyncWriteExt ;
15+ use tokio:: net:: TcpListener ;
16+ use tokio:: task:: JoinHandle ;
1717
1818const INFERENCE_PROVIDER_NAME : & str = "e2e-host-inference" ;
1919const INFERENCE_PROVIDER_UNREACHABLE_NAME : & str = "e2e-host-inference-unreachable" ;
20- const TEST_SERVER_IMAGE : & str = "public.ecr.aws/docker/library/python:3.13-alpine" ;
2120static INFERENCE_ROUTE_LOCK : Mutex < ( ) > = Mutex :: new ( ( ) ) ;
2221
2322async fn run_cli ( args : & [ & str ] ) -> Result < String , String > {
@@ -44,117 +43,63 @@ async fn run_cli(args: &[&str]) -> Result<String, String> {
4443 Ok ( combined)
4544}
4645
47- struct DockerServer {
46+ struct HostServer {
4847 port : u16 ,
49- container_id : String ,
48+ task : JoinHandle < ( ) > ,
5049}
5150
52- impl DockerServer {
51+ impl HostServer {
5352 async fn start ( response_body : & str ) -> Result < Self , String > {
54- let port = find_free_port ( ) ;
55- let script = r#"from http.server import BaseHTTPRequestHandler, HTTPServer
56- import os
57-
58- BODY = os.environ["RESPONSE_BODY"].encode()
59-
60- class Handler(BaseHTTPRequestHandler):
61- def do_GET(self):
62- self.send_response(200)
63- self.send_header("Content-Type", "application/json")
64- self.send_header("Content-Length", str(len(BODY)))
65- self.end_headers()
66- self.wfile.write(BODY)
67-
68- def do_POST(self):
69- length = int(self.headers.get("Content-Length", "0"))
70- if length:
71- self.rfile.read(length)
72- self.send_response(200)
73- self.send_header("Content-Type", "application/json")
74- self.send_header("Content-Length", str(len(BODY)))
75- self.end_headers()
76- self.wfile.write(BODY)
77-
78- def log_message(self, format, *args):
79- pass
80-
81- HTTPServer(("0.0.0.0", 8000), Handler).serve_forever()
82- "# ;
83-
84- let output = Command :: new ( "docker" )
85- . args ( [
86- "run" ,
87- "--detach" ,
88- "--rm" ,
89- "-e" ,
90- & format ! ( "RESPONSE_BODY={response_body}" ) ,
91- "-p" ,
92- & format ! ( "{port}:8000" ) ,
93- TEST_SERVER_IMAGE ,
94- "python3" ,
95- "-c" ,
96- script,
97- ] )
98- . output ( )
99- . map_err ( |e| format ! ( "start docker test server: {e}" ) ) ?;
100-
101- let stdout = String :: from_utf8_lossy ( & output. stdout ) . trim ( ) . to_string ( ) ;
102- let stderr = String :: from_utf8_lossy ( & output. stderr ) . to_string ( ) ;
103-
104- if !output. status . success ( ) {
105- return Err ( format ! (
106- "docker run failed (exit {:?}):\n {stderr}" ,
107- output. status. code( )
108- ) ) ;
109- }
110-
111- let server = Self {
112- port,
113- container_id : stdout,
114- } ;
115- server. wait_until_ready ( ) . await ?;
116- Ok ( server)
117- }
118-
119- async fn wait_until_ready ( & self ) -> Result < ( ) , String > {
120- let container_id = self . container_id . clone ( ) ;
121- timeout ( Duration :: from_secs ( 60 ) , async move {
122- let mut tick = interval ( Duration :: from_millis ( 500 ) ) ;
53+ let listener = TcpListener :: bind ( ( "0.0.0.0" , 0 ) )
54+ . await
55+ . map_err ( |e| format ! ( "bind host test server: {e}" ) ) ?;
56+ let port = listener
57+ . local_addr ( )
58+ . map_err ( |e| format ! ( "read host test server address: {e}" ) ) ?
59+ . port ( ) ;
60+ let response_body = response_body. as_bytes ( ) . to_vec ( ) ;
61+ let task = tokio:: spawn ( async move {
12362 loop {
124- tick. tick ( ) . await ;
125- let output = Command :: new ( "docker" )
126- . args ( [
127- "exec" ,
128- & container_id,
129- "python3" ,
130- "-c" ,
131- "import urllib.request; urllib.request.urlopen('http://127.0.0.1:8000', timeout=1).read()" ,
132- ] )
133- . output ( ) ;
134-
135- match output {
136- Ok ( result) if result. status . success ( ) => return Ok ( ( ) ) ,
137- Ok ( _) | Err ( _) => continue ,
138- }
63+ let Ok ( ( mut stream, _) ) = listener. accept ( ) . await else {
64+ break ;
65+ } ;
66+ let body = response_body. clone ( ) ;
67+ tokio:: spawn ( async move {
68+ let mut request = Vec :: new ( ) ;
69+ let mut buf = [ 0_u8 ; 1024 ] ;
70+ loop {
71+ let Ok ( read) = stream. read ( & mut buf) . await else {
72+ return ;
73+ } ;
74+ if read == 0 {
75+ return ;
76+ }
77+ request. extend_from_slice ( & buf[ ..read] ) ;
78+ if request. windows ( 4 ) . any ( |window| window == b"\r \n \r \n " ) {
79+ break ;
80+ }
81+ }
82+
83+ let response = format ! (
84+ "HTTP/1.1 200 OK\r \n Content-Type: application/json\r \n Content-Length: {}\r \n Connection: close\r \n \r \n " ,
85+ body. len( )
86+ ) ;
87+ if stream. write_all ( response. as_bytes ( ) ) . await . is_err ( ) {
88+ return ;
89+ }
90+ let _ = stream. write_all ( & body) . await ;
91+ let _ = stream. shutdown ( ) . await ;
92+ } ) ;
13993 }
140- } )
141- . await
142- . map_err ( |_| {
143- format ! (
144- "docker test server {} did not become ready within 60s" ,
145- self . container_id
146- )
147- } ) ?
94+ } ) ;
95+
96+ Ok ( Self { port, task } )
14897 }
14998}
15099
151- impl Drop for DockerServer {
100+ impl Drop for HostServer {
152101 fn drop ( & mut self ) {
153- let _ = Command :: new ( "docker" )
154- . args ( [ "rm" , "-f" , & self . container_id ] )
155- . stdout ( Stdio :: null ( ) )
156- . stderr ( Stdio :: null ( ) )
157- . status ( ) ;
102+ self . task . abort ( ) ;
158103 }
159104}
160105
@@ -245,7 +190,7 @@ network_policies:
245190
246191#[ tokio:: test]
247192async fn sandbox_reaches_host_openshell_internal_via_host_gateway_alias ( ) {
248- let server = DockerServer :: start ( r#"{"message":"hello-from-host"}"# )
193+ let server = HostServer :: start ( r#"{"message":"hello-from-host"}"# )
249194 . await
250195 . expect ( "start host echo server" ) ;
251196 let policy = write_policy ( server. port ) . expect ( "write custom policy" ) ;
@@ -292,7 +237,7 @@ async fn sandbox_inference_local_routes_to_host_openshell_internal() {
292237 return ;
293238 }
294239
295- let server = DockerServer :: start (
240+ let server = HostServer :: start (
296241 r#"{"id":"chatcmpl-test","object":"chat.completion","created":1,"model":"host-echo","choices":[{"index":0,"message":{"role":"assistant","content":"hello-from-host"},"finish_reason":"stop"}]}"# ,
297242 )
298243 . await
0 commit comments