| 452 | |
| 453 | #[tokio::test] |
| 454 | async fn full_bootstrap_join_flow() { |
| 455 | // Node 1 bootstraps, Node 2 joins via QUIC. |
| 456 | use crate::transport::credentials::TransportCredentials; |
| 457 | let t1 = Arc::new( |
| 458 | NexarTransport::new( |
| 459 | 1, |
| 460 | "127.0.0.1:0".parse().unwrap(), |
| 461 | TransportCredentials::Insecure, |
| 462 | ) |
| 463 | .unwrap(), |
| 464 | ); |
| 465 | let t2 = Arc::new( |
| 466 | NexarTransport::new( |
| 467 | 2, |
| 468 | "127.0.0.1:0".parse().unwrap(), |
| 469 | TransportCredentials::Insecure, |
| 470 | ) |
| 471 | .unwrap(), |
| 472 | ); |
| 473 | |
| 474 | let (_dir1, catalog1) = temp_catalog(); |
| 475 | let (_dir2, catalog2) = temp_catalog(); |
| 476 | |
| 477 | let addr1 = t1.local_addr(); |
| 478 | let addr2 = t2.local_addr(); |
| 479 | |
| 480 | let config1 = ClusterConfig { |
| 481 | node_id: 1, |
| 482 | listen_addr: addr1, |
| 483 | seed_nodes: vec![addr1], |
| 484 | num_groups: 2, |
| 485 | replication_factor: 1, |
| 486 | data_dir: _dir1.path().to_path_buf(), |
| 487 | force_bootstrap: false, |
| 488 | join_retry: Default::default(), |
| 489 | swim_udp_addr: None, |
| 490 | election_timeout_min: Duration::from_millis(150), |
| 491 | election_timeout_max: Duration::from_millis(300), |
| 492 | install_snapshot_chunk_bytes: 4 * 1024 * 1024, |
| 493 | orphan_partial_max_age_secs: 300, |
| 494 | }; |
| 495 | let state1 = bootstrap(&config1, &catalog1, None).unwrap(); |
| 496 | |
| 497 | // state1.topology and state1.routing are Arc<RwLock<T>> after the |
| 498 | // ClusterState refactor. |
| 499 | let topology1 = state1.topology.clone(); |
| 500 | let routing1 = state1.routing.clone(); |
| 501 | |
| 502 | struct JoinHandler { |
| 503 | topology: std::sync::Arc<std::sync::RwLock<ClusterTopology>>, |
| 504 | routing: std::sync::Arc<std::sync::RwLock<RoutingTable>>, |
| 505 | } |
| 506 | |
| 507 | impl crate::transport::RaftRpcHandler for JoinHandler { |
| 508 | async fn handle_rpc(&self, rpc: RaftRpc) -> Result<RaftRpc> { |
| 509 | match rpc { |
| 510 | RaftRpc::JoinRequest(req) => { |
| 511 | let mut topo = self.topology.write().unwrap(); |