(
path: &str,
id: u64,
)
| 423 | |
| 424 | #[cfg(target_family = "unix")] |
| 425 | async fn connect_scheduler( |
| 426 | path: &str, |
| 427 | id: u64, |
| 428 | ) -> ( |
| 429 | mpsc::UnboundedSender<WorkerMessage>, |
| 430 | mpsc::UnboundedReceiver<SchedulerMessage>, |
| 431 | ) { |
| 432 | let mut stream = tokio::net::UnixStream::connect(path) |
| 433 | .await |
| 434 | .expect("scheduler should have opened socket"); |
| 435 | |
| 436 | simpleproto::write_message(&WorkerMessage::Hello(id), &mut stream) |
| 437 | .await |
| 438 | .expect("should write to scheduler successfully"); |
| 439 | |
| 440 | let (mut reader_half, mut writer_half) = stream.into_split(); |
| 441 | |
| 442 | let scheduler_rx = { |
| 443 | let (tx, rx) = mpsc::unbounded_channel::<SchedulerMessage>(); |
| 444 | |
| 445 | tokio::spawn(async move { simpleproto::message_reader(&mut reader_half, tx).await }); |
| 446 | rx |
| 447 | }; |
| 448 | |
| 449 | let scheduler_tx = { |
| 450 | let (tx, rx) = mpsc::unbounded_channel::<WorkerMessage>(); |
| 451 | tokio::spawn(async move { simpleproto::message_writer(&mut writer_half, rx).await }); |
| 452 | |
| 453 | tx |
| 454 | }; |
| 455 | |
| 456 | (scheduler_tx, scheduler_rx) |
| 457 | } |
| 458 | |
| 459 | #[cfg(target_family = "windows")] |
| 460 | async fn connect_scheduler( |
no test coverage detected