Simulated TPC event loop using raw epoll + eventfd. This is what Glommio/monoio would do internally — poll an fd, wake up, process work, signal the other side.
(
mut data: nodedb_bridge::async_bridge::DataHandle<u64, u64>,
done: Arc<AtomicBool>,
)
| 28 | /// This is what Glommio/monoio would do internally — poll an fd, wake up, |
| 29 | /// process work, signal the other side. |
| 30 | fn tpc_event_loop( |
| 31 | mut data: nodedb_bridge::async_bridge::DataHandle<u64, u64>, |
| 32 | done: Arc<AtomicBool>, |
| 33 | ) { |
| 34 | let epoll_fd = unsafe { libc::epoll_create1(0) }; |
| 35 | assert!(epoll_fd >= 0, "epoll_create1 failed"); |
| 36 | |
| 37 | // Register the request-available eventfd with epoll (edge-triggered). |
| 38 | let req_fd = data.request_wake_fd(); |
| 39 | let mut event = libc::epoll_event { |
| 40 | events: (libc::EPOLLIN | libc::EPOLLET) as u32, |
| 41 | u64: req_fd as u64, |
| 42 | }; |
| 43 | let ret = unsafe { libc::epoll_ctl(epoll_fd, libc::EPOLL_CTL_ADD, req_fd, &mut event) }; |
| 44 | assert_eq!(ret, 0, "epoll_ctl failed"); |
| 45 | |
| 46 | let mut processed = 0u64; |
| 47 | let mut events = [libc::epoll_event { events: 0, u64: 0 }; 8]; |
| 48 | |
| 49 | while processed < MESSAGE_COUNT { |
| 50 | // epoll_wait with short timeout so we don't hang on edge-triggered misses. |
| 51 | let nfds = unsafe { libc::epoll_wait(epoll_fd, events.as_mut_ptr(), 8, 10) }; |
| 52 | |
| 53 | if nfds < 0 { |
| 54 | let err = std::io::Error::last_os_error(); |
| 55 | if err.kind() == std::io::ErrorKind::Interrupted { |
| 56 | continue; |
| 57 | } |
| 58 | panic!("epoll_wait failed: {err}"); |
| 59 | } |
| 60 | |
| 61 | // Consume the eventfd signal to re-arm edge trigger. |
| 62 | if nfds > 0 { |
| 63 | let _ = data.req_wake.consumer_wake.try_read(); |
| 64 | } |
| 65 | |
| 66 | // Drain all available requests. |
| 67 | let mut batch = Vec::new(); |
| 68 | data.drain_requests(&mut batch, 512); |
| 69 | |
| 70 | for req in batch { |
| 71 | loop { |
| 72 | match data.try_send_response(req * 2) { |
| 73 | Ok(()) => break, |
| 74 | Err(nodedb_bridge::BridgeError::Full { .. }) => { |
| 75 | thread::yield_now(); |
| 76 | } |
| 77 | Err(e) => panic!("TPC send error: {e}"), |
| 78 | } |
| 79 | } |
| 80 | processed += 1; |
| 81 | } |
| 82 | } |
| 83 | |
| 84 | done.store(true, Ordering::Release); |
| 85 | unsafe { libc::close(epoll_fd) }; |
| 86 | } |
| 87 |
no test coverage detected