(&mut self, req: Message)
| 77 | } |
| 78 | |
| 79 | async fn forward_request(&mut self, req: Message) -> Result<Tensor> { |
| 80 | // Only measure timing when debug logging is active — avoids Instant::now() overhead. |
| 81 | let send_start = log::log_enabled!(log::Level::Debug) |
| 82 | .then(std::time::Instant::now); |
| 83 | |
| 84 | req.to_writer_buf(&mut self.stream, &mut self.write_buf) |
| 85 | .await |
| 86 | .map_err(|e| anyhow!("error sending message {:?}: {}", req, e))?; |
| 87 | let send_elapsed = send_start.map(|s| s.elapsed()); |
| 88 | |
| 89 | let recv_start = log::log_enabled!(log::Level::Debug) |
| 90 | .then(std::time::Instant::now); |
| 91 | |
| 92 | let (resp_size, msg) = Message::from_reader_buf(&mut self.stream, &mut self.read_buf) |
| 93 | .await |
| 94 | .map_err(|e| anyhow!("error receiving response for {:?}: {}", req, e))?; |
| 95 | |
| 96 | if let (Some(se), Some(rs)) = (send_elapsed, recv_start) { |
| 97 | log::debug!( |
| 98 | " {} send={:.1}ms recv={:.1}ms ({})", |
| 99 | &self.address, |
| 100 | se.as_secs_f64() * 1000.0, |
| 101 | rs.elapsed().as_secs_f64() * 1000.0, |
| 102 | human_bytes::human_bytes(resp_size as f64), |
| 103 | ); |
| 104 | } |
| 105 | |
| 106 | match msg { |
| 107 | Message::Tensor(raw) => Ok(raw.into_tensor(&self.device)?), |
| 108 | Message::WorkerError { message } => Err(anyhow!( |
| 109 | "worker {} reported error: {}", |
| 110 | &self.address, |
| 111 | message |
| 112 | )), |
| 113 | _ => Err(anyhow!("unexpected response {:?}", &msg)), |
| 114 | } |
| 115 | } |
| 116 | } |
| 117 | |
| 118 | impl std::fmt::Display for Client { |
no test coverage detected