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