(
mut inbound: tonic::Streaming<TcpForwardFrame>,
relay_stream: tokio::io::DuplexStream,
tx: mpsc::Sender<Result<TcpForwardFrame, Status>>,
sandbox_id: &str,
channel_id: &str,
)
| 1147 | } |
| 1148 | |
| 1149 | async fn bridge_forward_tcp_stream( |
| 1150 | mut inbound: tonic::Streaming<TcpForwardFrame>, |
| 1151 | relay_stream: tokio::io::DuplexStream, |
| 1152 | tx: mpsc::Sender<Result<TcpForwardFrame, Status>>, |
| 1153 | sandbox_id: &str, |
| 1154 | channel_id: &str, |
| 1155 | ) { |
| 1156 | let (mut relay_read, mut relay_write) = tokio::io::split(relay_stream); |
| 1157 | |
| 1158 | let sandbox_id_in = sandbox_id.to_string(); |
| 1159 | let channel_id_in = channel_id.to_string(); |
| 1160 | tokio::spawn(async move { |
| 1161 | loop { |
| 1162 | match inbound.message().await { |
| 1163 | Ok(Some(frame)) => { |
| 1164 | let Some(openshell_core::proto::tcp_forward_frame::Payload::Data(data)) = |
| 1165 | frame.payload |
| 1166 | else { |
| 1167 | warn!(sandbox_id = %sandbox_id_in, channel_id = %channel_id_in, "ForwardTcp: received non-data frame after init"); |
| 1168 | break; |
| 1169 | }; |
| 1170 | if data.is_empty() { |
| 1171 | continue; |
| 1172 | } |
| 1173 | if let Err(err) = |
| 1174 | tokio::io::AsyncWriteExt::write_all(&mut relay_write, &data).await |
| 1175 | { |
| 1176 | warn!(sandbox_id = %sandbox_id_in, channel_id = %channel_id_in, error = %err, "ForwardTcp: write to relay failed"); |
| 1177 | break; |
| 1178 | } |
| 1179 | } |
| 1180 | Ok(None) => break, |
| 1181 | Err(err) => { |
| 1182 | debug!(sandbox_id = %sandbox_id_in, channel_id = %channel_id_in, error = %err, "ForwardTcp: inbound stream ended"); |
| 1183 | break; |
| 1184 | } |
| 1185 | } |
| 1186 | } |
| 1187 | let _ = tokio::io::AsyncWriteExt::shutdown(&mut relay_write).await; |
| 1188 | }); |
| 1189 | |
| 1190 | let mut buf = vec![0u8; TCP_FORWARD_CHUNK_SIZE]; |
| 1191 | loop { |
| 1192 | match tokio::io::AsyncReadExt::read(&mut relay_read, &mut buf).await { |
| 1193 | Ok(0) => break, |
| 1194 | Ok(n) => { |
| 1195 | let frame = TcpForwardFrame { |
| 1196 | payload: Some(openshell_core::proto::tcp_forward_frame::Payload::Data( |
| 1197 | buf[..n].to_vec(), |
| 1198 | )), |
| 1199 | }; |
| 1200 | if tx.send(Ok(frame)).await.is_err() { |
| 1201 | break; |
| 1202 | } |
| 1203 | } |
| 1204 | Err(err) => { |
| 1205 | warn!(sandbox_id = %sandbox_id, channel_id = %channel_id, error = %err, "ForwardTcp: read from relay failed"); |
| 1206 | let _ = tx |
no test coverage detected