From 50bb59449609787e0a19f4a8b8ee8c20e6f61240 Mon Sep 17 00:00:00 2001 From: 80347547 Date: Thu, 26 Mar 2026 17:01:31 +0800 Subject: [PATCH] fix(p2p-11): add reader and listener lifecycle hardening --- curvine-client/src/block/block_client.rs | 12 + curvine-client/src/block/block_client_pool.rs | 18 ++ curvine-client/src/block/block_reader.rs | 125 +++++++++-- .../src/block/block_reader_remote.rs | 91 +++++++- curvine-client/src/file/fs_context.rs | 82 ++++++- curvine-client/src/p2p/service.rs | 210 +++++++++++++++--- 6 files changed, 481 insertions(+), 57 deletions(-) diff --git a/curvine-client/src/block/block_client.rs b/curvine-client/src/block/block_client.rs index e20d84923..aee227f03 100644 --- a/curvine-client/src/block/block_client.rs +++ b/curvine-client/src/block/block_client.rs @@ -60,6 +60,18 @@ impl BlockClient { } } + #[cfg(test)] + pub(crate) fn new_for_test(worker_addr: WorkerAddress) -> Self { + Self { + client: None, + client_name: String::new(), + timeout: Duration::from_millis(0), + pool: None, + worker_addr, + uptime: LocalTime::mills(), + } + } + pub fn set_pool(&mut self, pool: Arc) { self.pool.replace(pool); self.uptime = LocalTime::mills(); diff --git a/curvine-client/src/block/block_client_pool.rs b/curvine-client/src/block/block_client_pool.rs index 9a29c6d90..ac6d4503d 100644 --- a/curvine-client/src/block/block_client_pool.rs +++ b/curvine-client/src/block/block_client_pool.rs @@ -215,4 +215,22 @@ impl BlockClientPool { .block_idle_conn .set(idle_count as i64); } + + pub fn shutdown(self: &Arc) { + let mut total_cleared = 0; + self.pool.retain(|_addr, queue| { + while let Some(mut client) = queue.pop_front() { + client.clear_pool(); + total_cleared += 1; + } + false + }); + self.cur_idle_size.set(0); + if total_cleared > 0 { + info!( + "shutdown cleared {} idle block connections from pool {}", + total_cleared, self.id + ); + } + } } diff --git a/curvine-client/src/block/block_reader.rs b/curvine-client/src/block/block_reader.rs index 0f7eeaea4..b7ad16dbc 100644 --- a/curvine-client/src/block/block_reader.rs +++ b/curvine-client/src/block/block_reader.rs @@ -17,6 +17,7 @@ use crate::block::{BlockReaderHole, BlockReaderLocal, BlockReaderRemote}; use crate::file::{FsContext, ReadChunkKey}; use crate::p2p::ChunkId; use bytes::Bytes; +use curvine_common::error::FsError; use curvine_common::state::{ClientAddress, ExtendedBlock, LocatedBlock, WorkerAddress}; use curvine_common::FsResult; use log::warn; @@ -60,6 +61,13 @@ impl ReaderAdapter { } } + async fn abort(&mut self) { + match self { + Local(_) | Hole(_) => {} + Remote(r) => r.abort().await, + } + } + fn remaining(&self) -> i64 { match self { Local(r) => r.remaining(), @@ -391,36 +399,109 @@ impl BlockReader { return Ok(chunk); } Err(e) => { - if matches!(&self.inner, Hole(_)) || self.locs.is_empty() { - return Err(e.ctx(format!( - "failed to read block on {}", - self.inner.worker_address() - ))); - } - - warn!( - "read data error block id {}, addr {}: {}", - self.block_id(), - self.inner.worker_address(), - e - ); - self.locs.retain(|x| x != self.inner.worker_address()); - self.inner = Self::get_reader( - &self.locs, - self.block.clone(), - self.fs_context.clone(), - self.pos(), - self.len(), - ) - .await?; + self.handle_worker_read_error(e).await?; } } } } + async fn handle_worker_read_error(&mut self, e: FsError) -> FsResult<()> { + if matches!(&self.inner, Hole(_)) || self.locs.is_empty() { + return Err(e.ctx(format!( + "failed to read block on {}", + self.inner.worker_address() + ))); + } + + let failed_addr = self.inner.worker_address().clone(); + warn!( + "read data error block id {}, addr {}: {}", + self.block_id(), + failed_addr, + e + ); + self.inner.abort().await; + self.locs.retain(|x| x != &failed_addr); + if self.locs.is_empty() { + return Err(e.ctx(format!("failed to read block on {}", failed_addr))); + } + self.inner = Self::get_reader( + &self.locs, + self.block.clone(), + self.fs_context.clone(), + self.pos(), + self.len(), + ) + .await?; + Ok(()) + } + fn advance_cached_position(&mut self, len: usize) -> FsResult<()> { let next_pos = self.pos() + len as i64; self.inner.seek(next_pos)?; Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::block::block_reader::ReaderAdapter::Remote; + use curvine_common::conf::ClusterConf; + use curvine_common::error::FsError; + use curvine_common::state::{ + FileAllocMode, FileAllocOpts, FileType, StorageType, WorkerAddress, + }; + use once_cell::sync::Lazy; + + static TEST_RT: Lazy> = Lazy::new(|| { + let conf = ClusterConf::default(); + Arc::new(conf.client_rpc_conf().create_runtime()) + }); + + fn test_fs_context() -> Arc { + let conf = ClusterConf::default(); + Arc::new(FsContext::with_rt(conf, TEST_RT.clone()).expect("fs context should build")) + } + + fn test_worker(worker_id: u32) -> WorkerAddress { + WorkerAddress { + worker_id, + hostname: format!("worker-{}", worker_id), + ip_addr: "127.0.0.1".to_string(), + rpc_port: 8000 + worker_id, + web_port: 9000 + worker_id, + } + } + + #[tokio::test] + async fn last_replica_read_failure_does_not_turn_allocated_block_into_hole() { + let fs_context = test_fs_context(); + let worker = test_worker(1); + let block = ExtendedBlock::with_alloc( + 7, + 4, + StorageType::Disk, + FileType::File, + Some(FileAllocOpts::with_alloc(4, FileAllocMode::ZERO_RANGE)), + ); + let remote = BlockReaderRemote::new_for_test(block.clone(), worker.clone(), 0, 4); + let mut reader = BlockReader { + inner: Remote(remote), + locs: vec![worker], + block, + file_id: 11, + file_version_epoch: 3, + file_mtime: 17, + fs_context, + }; + + let err = reader + .handle_worker_read_error(FsError::common("boom")) + .await + .expect_err("last replica failure should surface as error"); + + assert!(matches!(reader.inner, Remote(_))); + assert!(err.to_string().contains("failed to read block")); + } +} diff --git a/curvine-client/src/block/block_reader_remote.rs b/curvine-client/src/block/block_reader_remote.rs index ca4e0bef2..85f48882e 100644 --- a/curvine-client/src/block/block_reader_remote.rs +++ b/curvine-client/src/block/block_reader_remote.rs @@ -22,7 +22,7 @@ use orpc::err_box; use orpc::sys::DataSlice; pub struct BlockReaderRemote { - client: BlockClient, + client: Option, block: ExtendedBlock, worker_address: WorkerAddress, pos: i64, @@ -57,7 +57,7 @@ impl BlockReaderRemote { .await?; let reader = Self { - client, + client: Some(client), block, worker_address, pos: off, @@ -70,6 +70,25 @@ impl BlockReaderRemote { Ok(reader) } + #[cfg(test)] + pub(crate) fn new_for_test( + block: ExtendedBlock, + worker_address: WorkerAddress, + off: i64, + len: i64, + ) -> Self { + Self { + client: None, + block, + worker_address, + pos: off, + len, + req_id: Utils::req_id(), + seq_id: 0, + header: None, + } + } + fn next_seq_id(&mut self) -> i32 { self.seq_id += 1; self.seq_id @@ -108,7 +127,10 @@ impl BlockReaderRemote { let seq_id = self.next_seq_id(); let header = self.header.take(); - let chunk = self.client.read_data(self.req_id, seq_id, header).await?; + let Some(client) = self.client.as_ref() else { + return err_box!("No readable data"); + }; + let chunk = client.read_data(self.req_id, seq_id, header).await?; self.pos += chunk.len() as i64; Ok(chunk) @@ -116,11 +138,31 @@ impl BlockReaderRemote { pub async fn complete(&mut self) -> FsResult<()> { let next_seq_id = self.next_seq_id(); - self.client + let Some(client) = self.client.as_ref() else { + return Ok(()); + }; + client .read_commit(&self.block, self.req_id, next_seq_id) .await } + pub async fn abort(&mut self) { + let next_seq_id = self.next_seq_id(); + let Some(mut client) = self.client.take() else { + return; + }; + let _ = client + .read_commit(&self.block, self.req_id, next_seq_id) + .await; + client.clear_pool(); + self.header = None; + } + + #[cfg(test)] + pub(crate) fn set_test_client(&mut self, client: BlockClient) { + self.client = Some(client); + } + pub fn block_id(&self) -> i64 { self.block.id } @@ -129,3 +171,44 @@ impl BlockReaderRemote { &self.worker_address } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::block::BlockClientPool; + use curvine_common::state::{FileType, StorageType}; + use std::sync::Arc; + + fn test_worker(worker_id: u32) -> WorkerAddress { + WorkerAddress { + worker_id, + hostname: format!("worker-{}", worker_id), + ip_addr: "127.0.0.1".to_string(), + rpc_port: 8000 + worker_id, + web_port: 9000 + worker_id, + } + } + + #[tokio::test] + async fn aborted_remote_reader_does_not_repool_and_complete_is_noop() { + let worker = test_worker(1); + let block = ExtendedBlock::new(9, 4, StorageType::Disk, FileType::File); + let mut reader = BlockReaderRemote::new_for_test(block, worker.clone(), 0, 4); + let pool = Arc::new(BlockClientPool::new(true, 8, 60_000)); + let mut client = BlockClient::new_for_test(worker); + client.set_pool(pool.clone()); + reader.set_test_client(client); + + assert_eq!(pool.idle_conn(), 0); + assert_eq!(Arc::strong_count(&pool), 2); + + reader.abort().await; + + assert_eq!(pool.idle_conn(), 0); + assert_eq!(Arc::strong_count(&pool), 1); + reader + .complete() + .await + .expect("aborted reader completion should be a no-op"); + } +} diff --git a/curvine-client/src/file/fs_context.rs b/curvine-client/src/file/fs_context.rs index c6909d429..c460bbdad 100644 --- a/curvine-client/src/file/fs_context.rs +++ b/curvine-client/src/file/fs_context.rs @@ -311,19 +311,26 @@ impl FsContext { pub fn start_clean_task(fs: CurvineFileSystem, pool: Arc) { let metric_report_enable = fs.conf().client.metric_report_enable; let interval = fs.conf().client.clean_task_interval; - let fs_context = fs.fs_context.clone(); let rt = fs.clone_runtime(); - let metrics_fs = fs.clone(); + let fs_context = Arc::downgrade(&fs.fs_context); + let metrics_context = fs_context.clone(); + let pool = Arc::downgrade(&pool); rt.spawn(async move { let mut interval = tokio::time::interval(interval); loop { interval.tick().await; + let Some(pool) = pool.upgrade() else { + break; + }; pool.clear_idle_conn(); if metric_report_enable { - if let Err(e) = metrics_fs.metrics_report().await { + let Some(fs_context) = metrics_context.upgrade() else { + break; + }; + if let Err(e) = metrics_report(&fs_context).await { warn!("metrics report: {}", e) } } @@ -331,15 +338,22 @@ impl FsContext { }); rt.spawn(async move { - if let Err(e) = sync_p2p_runtime_policy(&fs_context).await { + let Some(startup_fs_context) = fs_context.upgrade() else { + return; + }; + if let Err(e) = sync_p2p_runtime_policy(&startup_fs_context).await { warn!("sync p2p runtime policy: {}", e) } + drop(startup_fs_context); let mut interval = tokio::time::interval(interval); interval.tick().await; loop { interval.tick().await; - if let Err(e) = sync_p2p_runtime_policy(&fs_context).await { + let Some(runtime_fs_context) = fs_context.upgrade() else { + break; + }; + if let Err(e) = sync_p2p_runtime_policy(&runtime_fs_context).await { warn!("sync p2p runtime policy: {}", e) } } @@ -371,13 +385,34 @@ async fn sync_p2p_runtime_policy(fs_context: &Arc) -> FsResult<()> { } } +async fn metrics_report(fs_context: &Arc) -> FsResult<()> { + let metrics = ClientMetrics::encode()?; + FsClient::new(fs_context.clone()) + .metrics_report(metrics) + .await +} + +impl Drop for FsContext { + fn drop(&mut self) { + if let Some(service) = self.p2p_service.as_ref() { + service.stop(); + } + self.block_pool.shutdown(); + } +} + #[cfg(test)] mod tests { use super::{FsContext, ReadChunkKey}; + use crate::block::BlockClient; + use crate::file::CurvineFileSystem; use crate::p2p::{ChunkId, P2pState}; use bytes::Bytes; use curvine_common::conf::ClusterConf; - use std::sync::Arc; + use curvine_common::state::WorkerAddress; + use orpc::runtime::RpcRuntime; + use std::sync::{Arc, Weak}; + use std::time::Duration; #[test] fn fs_context_skips_p2p_service_when_disabled() { @@ -423,4 +458,39 @@ mod tests { let fetched = rt.block_on(service_b.fetch_chunk(chunk_id, data.len(), Some(99))); assert_eq!(fetched, Some(data)); } + + #[test] + fn clean_tasks_do_not_keep_fs_context_alive() { + let conf = ClusterConf::default(); + let rt = Arc::new(conf.client_rpc_conf().create_runtime()); + let fs = CurvineFileSystem::with_rt(conf, rt.clone()).expect("fs should build"); + let weak = Arc::downgrade(&fs.fs_context); + + drop(fs); + rt.block_on(async { + tokio::time::sleep(Duration::from_millis(20)).await; + }); + + assert!(weak.upgrade().is_none()); + } + + #[test] + fn drop_shutdown_clears_idle_block_pool_cycles() { + let conf = ClusterConf::default(); + let rt = Arc::new(conf.client_rpc_conf().create_runtime()); + let ctx = Arc::new(FsContext::with_rt(conf, rt).expect("fs context should build")); + let weak_pool: Weak<_> = Arc::downgrade(&ctx.block_pool); + + let mut client = BlockClient::new_for_test(WorkerAddress { + worker_id: 1, + ..WorkerAddress::default() + }); + client.set_pool(ctx.block_pool.clone()); + ctx.block_pool.release(client); + assert_eq!(ctx.block_pool.idle_conn(), 1); + + drop(ctx); + + assert!(weak_pool.upgrade().is_none()); + } } diff --git a/curvine-client/src/p2p/service.rs b/curvine-client/src/p2p/service.rs index 43867019a..0e2f8bf1a 100644 --- a/curvine-client/src/p2p/service.rs +++ b/curvine-client/src/p2p/service.rs @@ -47,7 +47,7 @@ use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use tokio::net::{TcpListener, TcpStream}; use tokio::sync::mpsc::error::TrySendError; -use tokio::sync::{mpsc, oneshot, Mutex as AsyncMutex, Notify, Semaphore}; +use tokio::sync::{mpsc, oneshot, watch, Mutex as AsyncMutex, Notify, Semaphore}; use tokio::time::{sleep, timeout}; use tokio_util::compat::{Compat, TokioAsyncReadCompatExt}; @@ -124,6 +124,8 @@ struct DataPlaneTicket { struct DataPlaneState { port: Arc, started: Arc, + generation: Arc, + shutdown_tx: Arc>, tickets: Arc>, next_ticket: Arc, ticket_ttl_ms: i64, @@ -132,9 +134,12 @@ struct DataPlaneState { impl DataPlaneState { fn new(conf: &ClientP2pConf) -> Self { + let (shutdown_tx, _) = watch::channel(false); Self { port: Arc::new(AtomicU64::new(0)), started: Arc::new(AtomicBool::new(false)), + generation: Arc::new(AtomicU64::new(0)), + shutdown_tx: Arc::new(shutdown_tx), tickets: Arc::new(FastDashMap::default()), next_ticket: Arc::new(AtomicU64::new(1)), ticket_ttl_ms: data_plane_ticket_ttl(conf).as_millis().max(1) as i64, @@ -220,6 +225,30 @@ impl DataPlaneState { Err(poisoned) => *poisoned.into_inner() = updated, } } + + fn subscribe_shutdown(&self) -> watch::Receiver { + self.shutdown_tx.subscribe() + } + + fn clear_shutdown(&self) { + if *self.shutdown_tx.borrow() { + let _ = self.shutdown_tx.send(false); + } + } + + fn request_shutdown(&self) { + if !*self.shutdown_tx.borrow() { + let _ = self.shutdown_tx.send(true); + } + } + + fn next_generation(&self) -> u64 { + self.generation.fetch_add(1, Ordering::AcqRel) + 1 + } + + fn generation(&self) -> u64 { + self.generation.load(Ordering::Acquire) + } } #[derive(Debug)] @@ -948,6 +977,7 @@ impl P2pService { pub fn start(&self) -> bool { if self.is_enabled() { self.stats(); + self.data_plane.clear_shutdown(); self.ensure_data_plane_started(); self.network_warmup_done.store(false, Ordering::Relaxed); self.ensure_network_started(); @@ -962,6 +992,8 @@ impl P2pService { self.state.store(P2pState::Stopped as u8, Ordering::Relaxed); self.network_warmup_done.store(false, Ordering::Relaxed); PEER_STATS.remove(&self.peer_id); + self.data_plane.request_shutdown(); + force_reset_data_plane_listener_state(&self.data_plane); self.data_plane.tickets.clear(); self.stop_network(); } @@ -1802,6 +1834,7 @@ impl P2pService { { return; } + let generation = self.data_plane.next_generation(); let data_plane = self.data_plane.clone(); let cache_manager = self.cache_manager.clone(); let pending_fetched = self.pending_fetched.clone(); @@ -1810,35 +1843,57 @@ impl P2pService { let conf = self.conf.clone(); let stats = self.stats(); let future = async move { - let listener = match TcpListener::bind(("0.0.0.0", 0)).await { - Ok(listener) => listener, - Err(e) => { - data_plane.started.store(false, Ordering::Relaxed); - warn!("failed to bind p2p data plane listener: {}", e); - return; - } - }; - if let Ok(addr) = listener.local_addr() { - data_plane.port.store(addr.port() as u64, Ordering::Relaxed); - } + let mut shutdown_rx = data_plane.subscribe_shutdown(); loop { - match listener.accept().await { - Ok((stream, _)) => handle_data_plane_connection( - stream, - data_plane.clone(), - cache_manager.clone(), - pending_fetched.clone(), - hot_published_chunks.clone(), - tenant_whitelist.clone(), - conf.clone(), - stats.clone(), - ), + let listener = match TcpListener::bind(("0.0.0.0", 0)).await { + Ok(listener) => listener, Err(e) => { - warn!("p2p data plane accept failed: {}", e); - break; + clear_data_plane_listener_port(&data_plane, generation); + warn!("failed to bind p2p data plane listener: {}", e); + tokio::select! { + _ = shutdown_rx.changed() => break, + _ = tokio::time::sleep(Duration::from_millis(200)) => continue, + } + } + }; + if let Ok(addr) = listener.local_addr() { + set_data_plane_listener_port(&data_plane, generation, addr.port()); + } + loop { + tokio::select! { + _ = shutdown_rx.changed() => { + reset_data_plane_listener_state(&data_plane, generation); + return; + } + accepted = listener.accept() => { + match accepted { + Ok((stream, _)) => handle_data_plane_connection( + stream, + data_plane.clone(), + cache_manager.clone(), + pending_fetched.clone(), + hot_published_chunks.clone(), + tenant_whitelist.clone(), + conf.clone(), + stats.clone(), + ), + Err(e) => { + clear_data_plane_listener_port(&data_plane, generation); + warn!("p2p data plane accept failed: {}", e); + tokio::select! { + _ = shutdown_rx.changed() => { + reset_data_plane_listener_state(&data_plane, generation); + return; + } + _ = tokio::time::sleep(Duration::from_millis(200)) => break, + } + } + } + } } } } + reset_data_plane_listener_state(&data_plane, generation); }; if let Some(runtime) = &self.runtime { runtime.spawn(future); @@ -3572,6 +3627,28 @@ fn persist_runtime_policy_version(conf: &ClientP2pConf, policy_version: u64) -> && fs::write(path, policy_version.to_string()).is_ok() } +fn set_data_plane_listener_port(data_plane: &DataPlaneState, generation: u64, port: u16) { + if data_plane.generation() == generation { + data_plane.port.store(port as u64, Ordering::Relaxed); + } +} + +fn clear_data_plane_listener_port(data_plane: &DataPlaneState, generation: u64) { + if data_plane.generation() == generation { + data_plane.port.store(0, Ordering::Relaxed); + } +} + +fn force_reset_data_plane_listener_state(data_plane: &DataPlaneState) { + data_plane.port.store(0, Ordering::Relaxed); + data_plane.started.store(false, Ordering::Relaxed); +} + +fn reset_data_plane_listener_state(data_plane: &DataPlaneState, generation: u64) { + if data_plane.generation() == generation { + force_reset_data_plane_listener_state(data_plane); + } +} fn parse_bootstrap_peers(peers: &[String]) -> Vec<(PeerId, Multiaddr)> { peers .iter() @@ -5467,6 +5544,30 @@ mod tests { TEST_RT.clone() } + async fn wait_for_data_plane_port(service: &P2pService, deadline: Duration) -> Option { + let started = Instant::now(); + while started.elapsed() < deadline { + let port = service.data_plane.port(); + if port > 0 && service.data_plane.started.load(Ordering::Relaxed) { + return Some(port); + } + sleep(Duration::from_millis(25)).await; + } + None + } + + async fn wait_for_data_plane_stopped(service: &P2pService, deadline: Duration) -> Option<()> { + let started = Instant::now(); + while started.elapsed() < deadline { + if service.data_plane.port() == 0 && !service.data_plane.started.load(Ordering::Relaxed) + { + return Some(()); + } + sleep(Duration::from_millis(25)).await; + } + None + } + #[test] fn empty_peer_whitelist_remains_open_after_bootstrap_merge() { let bootstrap_peer_ids = HashSet::from([new_peer_id()]); @@ -6263,6 +6364,65 @@ mod tests { service.stop(); } + #[test] + fn clear_data_plane_listener_port_keeps_task_started() { + let conf = test_conf("data-plane-reset"); + let data_plane = DataPlaneState::new(&conf); + let generation = data_plane.next_generation(); + data_plane.started.store(true, Ordering::Relaxed); + data_plane.port.store(31_002, Ordering::Relaxed); + + clear_data_plane_listener_port(&data_plane, generation); + + assert!(data_plane.started.load(Ordering::Relaxed)); + assert_eq!(data_plane.port(), 0); + } + #[test] + fn stale_generation_cannot_reset_new_listener_state() { + let conf = test_conf("data-plane-generation"); + let data_plane = DataPlaneState::new(&conf); + let stale_generation = data_plane.next_generation(); + let current_generation = data_plane.next_generation(); + data_plane.started.store(true, Ordering::Relaxed); + set_data_plane_listener_port(&data_plane, current_generation, 31_002); + + reset_data_plane_listener_state(&data_plane, stale_generation); + + assert!(data_plane.started.load(Ordering::Relaxed)); + assert_eq!(data_plane.port(), 31_002); + } + + #[tokio::test] + async fn data_plane_listener_stops_and_restarts_cleanly() { + let mut conf = test_conf("data-plane-lifecycle"); + conf.enable = true; + conf.enable_mdns = false; + conf.enable_dht = false; + + let service = P2pService::new_with_runtime(conf, Some(test_runtime())); + service.start(); + + let first_port = wait_for_data_plane_port(&service, Duration::from_secs(2)) + .await + .expect("data plane should start"); + assert!(first_port > 0); + + service.stop(); + wait_for_data_plane_stopped(&service, Duration::from_secs(2)) + .await + .expect("data plane should stop"); + + service.start(); + let restarted_port = wait_for_data_plane_port(&service, Duration::from_secs(2)) + .await + .expect("data plane should restart"); + assert!(restarted_port > 0); + + service.stop(); + wait_for_data_plane_stopped(&service, Duration::from_secs(2)) + .await + .expect("data plane should stop after restart"); + } #[tokio::test] async fn invalid_runtime_peer_whitelist_update_is_rejected() { let mut conf = test_conf("invalid-runtime-peer-whitelist");