MCPcopy Create free account
hub / github.com/NodeDB-Lab/nodedb / tpc_event_loop

Function tpc_event_loop

nodedb-bridge/tests/cross_runtime.rs:30–86  ·  view source on GitHub ↗

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>,
)

Source from the content-addressed store, hash-verified

28/// This is what Glommio/monoio would do internally — poll an fd, wake up,
29/// process work, signal the other side.
30fn 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

Callers 1

tokio_to_tpc_via_eventfdFunction · 0.85

Calls 7

request_wake_fdMethod · 0.80
kindMethod · 0.80
try_readMethod · 0.80
try_send_responseMethod · 0.80
as_mut_ptrMethod · 0.45
drain_requestsMethod · 0.45
storeMethod · 0.45

Tested by

no test coverage detected