()
| 14 | /// sends a query, and gets a response. |
| 15 | #[tokio::test] |
| 16 | async fn pgwire_connect_and_query() { |
| 17 | // Set up infrastructure (mirrors session::tests::full_request_response_roundtrip). |
| 18 | let dir = tempfile::tempdir().unwrap(); |
| 19 | let wal_path = dir.path().join("test.wal"); |
| 20 | let wal = Arc::new(WalManager::open_for_testing(&wal_path).unwrap()); |
| 21 | |
| 22 | let (dispatcher, data_sides) = Dispatcher::new(1, 64); |
| 23 | let shared = SharedState::new(dispatcher, wal); |
| 24 | |
| 25 | // Start a Data Plane core in a background thread. |
| 26 | let data_side = data_sides.into_iter().next().unwrap(); |
| 27 | let core_dir = dir.path().to_path_buf(); |
| 28 | let (core_stop_tx, core_stop_rx) = std::sync::mpsc::channel::<()>(); |
| 29 | let core_handle = tokio::task::spawn_blocking(move || { |
| 30 | let mut core = CoreLoop::open( |
| 31 | 0, |
| 32 | data_side.request_rx, |
| 33 | data_side.response_tx, |
| 34 | &core_dir, |
| 35 | std::sync::Arc::new(nodedb_types::OrdinalClock::new()), |
| 36 | ) |
| 37 | .unwrap(); |
| 38 | while matches!( |
| 39 | core_stop_rx.try_recv(), |
| 40 | Err(std::sync::mpsc::TryRecvError::Empty) |
| 41 | ) { |
| 42 | core.tick(); |
| 43 | std::thread::sleep(Duration::from_millis(1)); |
| 44 | } |
| 45 | }); |
| 46 | |
| 47 | // Start response poller. |
| 48 | let shared_poller = Arc::clone(&shared); |
| 49 | let (poller_shutdown_tx, mut poller_shutdown_rx) = tokio::sync::watch::channel(false); |
| 50 | let poller_handle = tokio::spawn(async move { |
| 51 | loop { |
| 52 | shared_poller.poll_and_route_responses(); |
| 53 | tokio::select! { |
| 54 | _ = tokio::time::sleep(Duration::from_millis(1)) => {} |
| 55 | _ = poller_shutdown_rx.changed() => break, |
| 56 | } |
| 57 | } |
| 58 | }); |
| 59 | |
| 60 | // Bind pgwire listener on random port. |
| 61 | let pg_listener = PgListener::bind("127.0.0.1:0".parse().unwrap()) |
| 62 | .await |
| 63 | .unwrap(); |
| 64 | let pg_addr = pg_listener.local_addr(); |
| 65 | |
| 66 | let (shutdown_bus, _) = |
| 67 | nodedb::control::shutdown::ShutdownBus::new(Arc::clone(&shared.shutdown)); |
| 68 | let shared_pg = Arc::clone(&shared); |
| 69 | let test_startup_gate = Arc::clone(&shared.startup); |
| 70 | let bus_pg = shutdown_bus.clone(); |
| 71 | let pg_handle = tokio::spawn(async move { |
| 72 | pg_listener |
| 73 | .run( |
nothing calls this directly
no test coverage detected