MCPcopy Create free account
hub / github.com/NVIDIA/OpenShell / handle_connection

Function handle_connection

crates/openshell-sandbox/src/metadata_server.rs:91–137  ·  view source on GitHub ↗
(
    handler: &H,
    mut stream: tokio::net::TcpStream,
)

Source from the content-addressed store, hash-verified

89const READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
90
91async fn handle_connection<H: MetadataHandler>(
92 handler: &H,
93 mut stream: tokio::net::TcpStream,
94) -> Result<()> {
95 let mut buf = vec![0u8; MAX_REQUEST_BYTES];
96 let mut used = 0;
97 let deadline = tokio::time::sleep(READ_TIMEOUT);
98 tokio::pin!(deadline);
99 loop {
100 tokio::select! {
101 result = stream.read(&mut buf[used..]) => {
102 let n = result.map_err(|e| miette::miette!("{e}"))?;
103 if n == 0 {
104 return Ok(());
105 }
106 used += n;
107 if buf[..used].windows(4).any(|w| w == b"\r\n\r\n") {
108 break;
109 }
110 if used >= buf.len() {
111 let _ = stream
112 .write_all(b"HTTP/1.1 413 Request Entity Too Large\r\nContent-Length: 0\r\n\r\n")
113 .await;
114 return Ok(());
115 }
116 }
117 () = &mut deadline => {
118 return Ok(());
119 }
120 }
121 }
122 let request = String::from_utf8_lossy(&buf[..used]);
123 let request_line = request.split("\r\n").next().unwrap_or("");
124 let mut parts = request_line.split_whitespace();
125 let method = parts.next().unwrap_or("");
126 let path = parts.next().unwrap_or("/");
127
128 tokio::time::timeout(
129 READ_TIMEOUT,
130 handler.handle(method, path, &buf[..used], &mut stream),
131 )
132 .await
133 .unwrap_or_else(|_| {
134 debug!(method, path, "metadata handler timed out");
135 Ok(())
136 })
137}

Callers 1

runFunction · 0.70

Calls 2

handleMethod · 0.80
nextMethod · 0.45

Tested by

no test coverage detected