()
| 31 | |
| 32 | impl NativeTestServer { |
| 33 | async fn start() -> Self { |
| 34 | let dir = tempfile::tempdir().unwrap(); |
| 35 | let wal_path = dir.path().join("test.wal"); |
| 36 | let wal = Arc::new(WalManager::open_for_testing(&wal_path).unwrap()); |
| 37 | |
| 38 | let (dispatcher, data_sides) = Dispatcher::new(1, 64); |
| 39 | let (event_producers, event_consumers) = create_event_bus(1); |
| 40 | |
| 41 | let shared = SharedState::new(dispatcher, Arc::clone(&wal)); |
| 42 | |
| 43 | let data_side = data_sides.into_iter().next().unwrap(); |
| 44 | let core_dir = dir.path().to_path_buf(); |
| 45 | let event_producer = event_producers.into_iter().next().unwrap(); |
| 46 | let core_array_catalog = shared.array_catalog.clone(); |
| 47 | let (core_stop_tx, core_stop_rx) = std::sync::mpsc::channel::<()>(); |
| 48 | let _core_handle = tokio::task::spawn_blocking(move || { |
| 49 | let mut core = CoreLoop::open_with_array_catalog( |
| 50 | 0, |
| 51 | data_side.request_rx, |
| 52 | data_side.response_tx, |
| 53 | &core_dir, |
| 54 | std::sync::Arc::new(nodedb_types::OrdinalClock::new()), |
| 55 | core_array_catalog, |
| 56 | ) |
| 57 | .unwrap(); |
| 58 | core.set_event_producer(event_producer); |
| 59 | while matches!( |
| 60 | core_stop_rx.try_recv(), |
| 61 | Err(std::sync::mpsc::TryRecvError::Empty) |
| 62 | ) { |
| 63 | core.tick(); |
| 64 | std::thread::sleep(Duration::from_millis(1)); |
| 65 | } |
| 66 | }); |
| 67 | |
| 68 | let shared_poller = Arc::clone(&shared); |
| 69 | let (poller_shutdown_tx, mut poller_shutdown_rx) = tokio::sync::watch::channel(false); |
| 70 | let _poller_handle = tokio::spawn(async move { |
| 71 | loop { |
| 72 | shared_poller.poll_and_route_responses(); |
| 73 | tokio::select! { |
| 74 | _ = tokio::time::sleep(Duration::from_millis(1)) => {} |
| 75 | _ = poller_shutdown_rx.changed() => break, |
| 76 | } |
| 77 | } |
| 78 | }); |
| 79 | |
| 80 | let watermark_store = |
| 81 | Arc::new(nodedb::event::watermark::WatermarkStore::open(dir.path()).unwrap()); |
| 82 | let trigger_dlq = Arc::new(std::sync::Mutex::new( |
| 83 | nodedb::event::trigger::TriggerDlq::open(dir.path()).unwrap(), |
| 84 | )); |
| 85 | let _event_plane = EventPlane::spawn( |
| 86 | event_consumers, |
| 87 | Arc::clone(&wal), |
| 88 | watermark_store, |
| 89 | Arc::clone(&shared), |
| 90 | trigger_dlq, |
nothing calls this directly
no test coverage detected