Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions curvine-client/src/block/block_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<BlockClientPool>) {
self.pool.replace(pool);
self.uptime = LocalTime::mills();
Expand Down
18 changes: 18 additions & 0 deletions curvine-client/src/block/block_client_pool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -215,4 +215,22 @@ impl BlockClientPool {
.block_idle_conn
.set(idle_count as i64);
}

pub fn shutdown(self: &Arc<Self>) {
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
);
}
}
}
125 changes: 103 additions & 22 deletions curvine-client/src/block/block_reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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<Arc<Runtime>> = Lazy::new(|| {
let conf = ClusterConf::default();
Arc::new(conf.client_rpc_conf().create_runtime())
});

fn test_fs_context() -> Arc<FsContext> {
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"));
}
}
91 changes: 87 additions & 4 deletions curvine-client/src/block/block_reader_remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ use orpc::err_box;
use orpc::sys::DataSlice;

pub struct BlockReaderRemote {
client: BlockClient,
client: Option<BlockClient>,
block: ExtendedBlock,
worker_address: WorkerAddress,
pos: i64,
Expand Down Expand Up @@ -57,7 +57,7 @@ impl BlockReaderRemote {
.await?;

let reader = Self {
client,
client: Some(client),
block,
worker_address,
pos: off,
Expand All @@ -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
Expand Down Expand Up @@ -108,19 +127,42 @@ 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)
}

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
}
Expand All @@ -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");
}
}
Loading
Loading