| 353 | } |
| 354 | |
| 355 | async fn _request<S, D>( |
| 356 | &self, |
| 357 | target_role: Option<Role>, |
| 358 | pl: &S, |
| 359 | ans: &mut D, |
| 360 | async_resp: Option<Receiver<Vec<u8>>>, |
| 361 | be_req_log: &mut stream::BackendInterfacesRequest, |
| 362 | ) -> Result<()> |
| 363 | where |
| 364 | S: ?Sized + serde::ser::Serialize + BasePayloadProvider, |
| 365 | D: serde::de::DeserializeOwned + BasePayloadResultProvider, |
| 366 | { |
| 367 | let server = if self.config.use_target_role_suffix { |
| 368 | match target_role { |
| 369 | Some(Role::FNS) => format!("{}/fns", self.config.server), |
| 370 | Some(Role::SNS) => format!("{}/sns", self.config.server), |
| 371 | Some(Role::HNS) => format!("{}/hns", self.config.server), |
| 372 | None => self.config.server.clone(), |
| 373 | } |
| 374 | } else { |
| 375 | self.config.server.clone() |
| 376 | }; |
| 377 | |
| 378 | let body = serde_json::to_string(&pl)?; |
| 379 | be_req_log.request_body.clone_from(&body); |
| 380 | |
| 381 | info!(server = %server, async_interface = %async_resp.is_some(), "Making request"); |
| 382 | |
| 383 | let res = self |
| 384 | .client |
| 385 | .post(&server) |
| 386 | .headers(self.headers.clone()) |
| 387 | .body(body) |
| 388 | .send() |
| 389 | .await? |
| 390 | .error_for_status()?; |
| 391 | |
| 392 | let resp_json = match async_resp { |
| 393 | Some(rx) => { |
| 394 | let sleep = tokio::time::sleep(self.config.async_timeout); |
| 395 | |
| 396 | tokio::select! { |
| 397 | rx_ans = rx => { |
| 398 | info!("Async response received"); |
| 399 | String::from_utf8(rx_ans?)? |
| 400 | } |
| 401 | _ = sleep => { |
| 402 | error!("Async response timeout"); |
| 403 | return Err(anyhow!("Async timeout")); |
| 404 | } |
| 405 | } |
| 406 | } |
| 407 | None => res.text().await?, |
| 408 | }; |
| 409 | |
| 410 | be_req_log.response_body.clone_from(&resp_json); |
| 411 | |
| 412 | let base: BasePayloadResult = serde_json::from_str(&resp_json)?; |