| 470 | } |
| 471 | |
| 472 | async fn request(&self, mut payload: Json) -> FlowResult<Json> { |
| 473 | let receiver = self.send_request(&mut payload)?; |
| 474 | let response_task = tokio::task::spawn_blocking(move || receiver.recv()); |
| 475 | let envelope = match tokio::time::timeout(WORKER_RPC_TIMEOUT, response_task).await { |
| 476 | Ok(result) => result |
| 477 | .map_err(|err| FlowError::Internal(format!("worker response task failed: {err}")))? |
| 478 | .map_err(|err| { |
| 479 | FlowError::Internal(format!("worker response channel closed: {err}")) |
| 480 | })?, |
| 481 | Err(_) => { |
| 482 | self.shutdown(); |
| 483 | return Err(FlowError::Internal(format!( |
| 484 | "worker request timed out after {} seconds", |
| 485 | WORKER_RPC_TIMEOUT.as_secs() |
| 486 | ))); |
| 487 | } |
| 488 | }; |
| 489 | worker_result(envelope) |
| 490 | } |
| 491 | |
| 492 | fn request_blocking(&self, mut payload: Json, timeout: Duration) -> FlowResult<WorkerEnvelope> { |
| 493 | let receiver = self.send_request(&mut payload)?; |