use std::sync::Arc;
use bytes::{Bytes, BytesMut};
use tokio::sync::RwLock;
use tracing::{debug, warn};
use crate::block::mapper::{BlockMapper, BlockReadPlan};
use crate::block::router::WorkerRouterView;
use crate::block::short_circuit::{ShortCircuitError, ShortCircuitFactory};
use crate::cache::{
page_cache_eligible, read_through_cache, CacheManager, ExternalRangeReader, FillMode,
};
use crate::client::worker::WorkerClientPool;
use crate::client::WorkerClient;
use crate::config::GoosefsConfig;
use crate::context::FileSystemContext;
use crate::error::{Error, Result};
use crate::io::reader::GrpcBlockReader;
use crate::proto::grpc::block::WorkerInfo;
use crate::proto::grpc::file::FileInfo;
use crate::proto::proto::dataserver::OpenUfsBlockOptions;
pub struct GoosefsFileReader {
config: GoosefsConfig,
path: String,
file_info: Arc<FileInfo>,
router: WorkerRouterView,
worker_pool: Option<Arc<WorkerClientPool>>,
_context: Option<Arc<FileSystemContext>>,
plans: Vec<BlockReadPlan>,
current_plan_index: usize,
total_bytes_read: u64,
offset: u64,
length: u64,
cache: Option<Arc<dyn CacheManager>>,
cache_file_id: Arc<str>,
cache_page_size: u64,
cache_fill: bool,
cache_async_write: bool,
short_circuit: Option<Arc<ShortCircuitFactory>>,
ufs_read_options: Option<OpenUfsBlockOptions>,
}
const ON_FILE_OPEN_DEDUP_TTL: std::time::Duration = std::time::Duration::from_secs(60);
type OnFileOpenDedupMap = std::collections::HashMap<(i64, i64, i64), std::time::Instant>;
static ON_FILE_OPEN_CACHE: std::sync::LazyLock<RwLock<OnFileOpenDedupMap>> =
std::sync::LazyLock::new(|| RwLock::new(std::collections::HashMap::new()));
impl GoosefsFileReader {
pub async fn open_with_context(ctx: Arc<FileSystemContext>, path: &str) -> Result<Self> {
let (file_info, router) = Self::init_with_context(&ctx, path).await?;
let file_length = file_info.length.unwrap_or(0) as u64;
let config = ctx.config().clone();
let pool = Some(ctx.acquire_worker_pool());
let mut reader = Self::build(
&config,
path,
file_info,
router,
pool,
Some(ctx.clone()),
0,
file_length,
)?;
reader.attach_cache(&ctx).await;
Ok(reader)
}
pub async fn open_range_with_context(
ctx: Arc<FileSystemContext>,
path: &str,
offset: u64,
length: u64,
) -> Result<Self> {
let (file_info, router) = Self::init_with_context(&ctx, path).await?;
let config = ctx.config().clone();
let pool = Some(ctx.acquire_worker_pool());
let mut reader = Self::build(
&config,
path,
file_info,
router,
pool,
Some(ctx.clone()),
offset,
length,
)?;
reader.attach_cache(&ctx).await;
Ok(reader)
}
async fn init_with_context(
ctx: &Arc<FileSystemContext>,
path: &str,
) -> Result<(Arc<FileInfo>, WorkerRouterView)> {
let file_info_cache = ctx.acquire_file_info_cache();
let file_info = if let Some(cached) = file_info_cache.as_ref().and_then(|c| c.get(path)) {
debug!(path = %path, "FileInfo cache hit (§A3 + S3 — Arc clone, zero deep copy)");
cached
} else {
let master = ctx.acquire_master();
let fetched = master.get_status(path).await?;
let arc_fetched = Arc::new(fetched);
if let Some(cache) = &file_info_cache {
cache.insert_arc(path, Arc::clone(&arc_fetched));
}
arc_fetched
};
let file_length = file_info.length.unwrap_or(0);
if file_length == 0 {
debug!(path = %path, "file is empty");
}
debug!(
path = %path,
file_length = file_length,
block_count = file_info.block_ids.len(),
block_size = ?file_info.block_size_bytes,
"fetched file metadata (via context)"
);
let shared_router = ctx.acquire_router();
if shared_router.workers_is_empty() {
return Err(Error::NoWorkerAvailable {
message: "no workers available for reading".to_string(),
});
}
let router = WorkerRouterView::from_shared(&shared_router);
debug!("reusing worker snapshot from context (A1)");
Ok((file_info, router))
}
#[allow(clippy::too_many_arguments)]
fn build(
config: &GoosefsConfig,
path: &str,
file_info: Arc<FileInfo>,
router: WorkerRouterView,
worker_pool: Option<Arc<WorkerClientPool>>,
context: Option<Arc<FileSystemContext>>,
offset: u64,
length: u64,
) -> Result<Self> {
let plans = BlockMapper::plan_read(&file_info, offset, length);
debug!(
path = %path,
offset = offset,
length = length,
block_segments = plans.len(),
"read plan created"
);
let ufs_read_options = {
let ufs_path = file_info.ufs_path.as_ref();
if ufs_path.map_or(true, |p| p.is_empty()) {
None
} else {
let block_size = file_info.block_size_bytes.unwrap_or(64 * 1024 * 1024);
Some(OpenUfsBlockOptions {
ufs_path: file_info.ufs_path.clone(),
offset_in_file: Some(0),
block_size: Some(block_size),
max_ufs_read_concurrency: None,
mount_id: file_info.mount_id,
no_cache: Some(!file_info.cacheable.unwrap_or(true)),
user: None,
caller_type: None,
})
}
};
Ok(Self {
config: config.clone(),
path: path.to_string(),
file_info,
router,
worker_pool,
_context: context,
plans,
current_plan_index: 0,
total_bytes_read: 0,
offset,
length,
cache: None,
cache_file_id: Arc::from(""),
cache_page_size: config.client_cache_page_size,
cache_fill: false,
cache_async_write: false,
short_circuit: None,
ufs_read_options,
})
}
async fn attach_cache(&mut self, ctx: &Arc<FileSystemContext>) {
self.short_circuit = ctx.acquire_short_circuit();
let file_id = self.file_info.file_id.unwrap_or(0);
if !page_cache_eligible(file_id) {
self.cache = None;
self.cache_fill = false;
return;
}
let cache = ctx.acquire_cache_manager();
let cache_file_id: Arc<str> = Arc::from(file_id.to_string());
if let Some(cache) = &cache {
let length = self.file_info.length.unwrap_or(0);
let mtime = self.file_info.last_modification_time_ms.unwrap_or(0);
let dedup_key = (file_id, length, mtime);
let skip = {
let guard = ON_FILE_OPEN_CACHE.read().await;
guard
.get(&dedup_key)
.is_some_and(|last| last.elapsed() < ON_FILE_OPEN_DEDUP_TTL)
};
if !skip {
cache.on_file_open(&cache_file_id, length, mtime).await;
let mut guard = ON_FILE_OPEN_CACHE.write().await;
if guard.len() > 4096 {
guard.retain(|_, last| last.elapsed() < ON_FILE_OPEN_DEDUP_TTL);
}
guard.insert(dedup_key, std::time::Instant::now());
}
}
let cfg = ctx.config();
self.cache_fill = cache.is_some(); self.cache_async_write = cfg.client_cache_async_write_enabled;
self.cache_page_size = cfg.client_cache_page_size;
self.cache_file_id = cache_file_id;
self.cache = cache;
}
pub async fn read_next_block(&mut self) -> Result<Option<Bytes>> {
let (abs_offset, abs_end) = loop {
if self.current_plan_index >= self.plans.len() {
return Ok(None);
}
let plan = self.plans[self.current_plan_index].clone();
let block_id = self.resolve_block_id(&plan);
if block_id <= 0 {
warn!(
block_index = plan.block_index,
plan_block_id = plan.block_id,
"invalid block ID, skipping block"
);
self.current_plan_index += 1;
continue;
}
let block_size = self.file_info.block_size_bytes.unwrap_or(64 * 1024 * 1024) as u64;
let abs_offset = (plan.block_index * block_size + plan.offset_in_block) as i64;
break (abs_offset, abs_offset + plan.length as i64);
};
let data = match self.cache.clone() {
Some(cache) => {
let file_id = self.cache_file_id.clone();
let page_size = self.cache_page_size;
let file_length = self.file_length() as i64;
let fill_mode = if !self.cache_fill {
FillMode::None
} else if self.cache_async_write {
FillMode::Async
} else {
FillMode::Sync
};
read_through_cache(
&cache,
self,
&file_id,
page_size,
file_length,
abs_offset,
abs_end,
fill_mode,
)
.await?
}
None => self.read_file_range(abs_offset, abs_end).await?,
};
let bytes_read = data.len() as u64;
self.total_bytes_read += bytes_read;
self.current_plan_index += 1;
debug!(
abs_offset = abs_offset,
abs_end = abs_end,
bytes_read = bytes_read,
total_read = self.total_bytes_read,
cache_enabled = self.cache.is_some(),
"block segment read complete"
);
Ok(Some(data))
}
async fn read_segment(&self, block_id: i64, plan: &BlockReadPlan) -> Result<Bytes> {
if let Some(sc_result) = self.try_short_circuit_read(block_id, plan).await {
return sc_result;
}
let ufs_options = self.build_ufs_read_options(plan);
let worker_info = self.router.select_worker(block_id).await?;
let worker_addr = Self::worker_addr(&worker_info)?;
let worker = match self.acquire_worker(&worker_addr).await {
Ok(w) => w,
Err(e) if e.is_authentication_failed() => {
debug!(
worker = %worker_addr,
error = %e,
"authentication failed on connect, reconnecting"
);
self.reconnect_worker(&worker_addr, None).await?
}
Err(e) => {
if let Some(a) = worker_info.address.as_ref() {
self.router.mark_failed(a);
}
warn!(
worker = %worker_addr,
error = %e,
"worker connection failed, trying another worker"
);
let retry = self.router.select_worker(block_id).await.map_err(|_| e)?;
let retry_addr = Self::worker_addr(&retry)?;
self.acquire_worker(&retry_addr).await?
}
};
let worker_generation = worker.generation();
match self
.try_read_block(&worker, block_id, plan, ufs_options.clone())
.await
{
Ok(d) => Ok(d),
Err(e) if e.is_authentication_failed() => {
debug!(
block_id = block_id,
worker = %worker_addr,
stale_generation = worker_generation,
error = %e,
"auth failed during block read, requesting single-flight reconnect"
);
let fresh = self
.reconnect_worker(&worker_addr, Some(worker_generation))
.await?;
self.try_read_block(&fresh, block_id, plan, ufs_options)
.await
}
Err(e) => Err(e),
}
}
async fn read_file_range(&self, abs_offset: i64, abs_end: i64) -> Result<Bytes> {
if abs_end <= abs_offset {
return Ok(Bytes::new());
}
let plans = BlockMapper::plan_read(
&self.file_info,
abs_offset as u64,
(abs_end - abs_offset) as u64,
);
let mut buf = BytesMut::with_capacity((abs_end - abs_offset) as usize);
for plan in &plans {
let block_id = self.resolve_block_id(plan);
if block_id <= 0 {
continue;
}
let data = self.read_segment(block_id, plan).await?;
if data.is_empty() {
return Err(Error::Internal {
message: format!("read_file_range: 0 bytes for block {block_id}"),
source: None,
});
}
buf.extend_from_slice(&data);
}
Ok(buf.freeze())
}
async fn try_short_circuit_read(
&self,
block_id: i64,
plan: &BlockReadPlan,
) -> Option<Result<Bytes>> {
let factory = match self.short_circuit.as_ref() {
Some(f) => f,
None => {
crate::metrics::counter(crate::metrics::name::CLIENT_SC_DECISION_SKIPPED).inc(1);
return None;
}
};
let block_size = self.block_logical_size(plan.block_index);
if !factory.should_use(block_id, block_size).await {
crate::metrics::counter(crate::metrics::name::CLIENT_SC_DECISION_SKIPPED).inc(1);
return None;
}
let reader = match factory.get_or_open(block_id, block_size).await {
Ok(r) => r,
Err(e) => {
debug!(
block_id = block_id,
error = %e,
"short-circuit open failed, falling back to gRPC"
);
crate::metrics::counter(crate::metrics::name::CLIENT_SC_DECISION_FALLBACK_OPEN)
.inc(1);
return None;
}
};
match reader.read_bytes(plan.offset_in_block as usize, plan.length as usize) {
Ok(bytes) => {
crate::metrics::counter(crate::metrics::name::CLIENT_SC_DECISION_HIT).inc(1);
Some(Ok(bytes))
}
Err(ShortCircuitError::OutOfRange {
off,
len,
file_size,
}) => {
crate::metrics::counter(crate::metrics::name::CLIENT_SC_DECISION_SEMANTIC_ERROR)
.inc(1);
Some(Err(Error::InvalidArgument {
message: format!(
"short-circuit read out of range on block {block_id}: \
off={off} len={len} block_size={file_size}"
),
}))
}
Err(e) => {
debug!(
block_id = block_id,
error = %e,
"short-circuit read failed, falling back to gRPC"
);
crate::metrics::counter(crate::metrics::name::CLIENT_SC_DECISION_FALLBACK_READ)
.inc(1);
factory.invalidate(block_id).await;
None
}
}
}
fn block_logical_size(&self, block_index: u64) -> i64 {
let bs = self.file_info.block_size_bytes.unwrap_or(64 * 1024 * 1024);
let file_length = self.file_length() as i64;
if bs <= 0 {
return file_length.max(0);
}
let start = block_index as i64 * bs;
(file_length - start).clamp(0, bs)
}
fn worker_addr(worker_info: &WorkerInfo) -> Result<String> {
let addr = worker_info
.address
.as_ref()
.ok_or_else(|| Error::Internal {
message: "worker has no address".to_string(),
source: None,
})?;
Ok(crate::block::router::rpc_endpoint(addr))
}
pub async fn read_all(&mut self) -> Result<Bytes> {
let expected_len = self.plans.iter().map(|p| p.length).sum::<u64>();
let mut buf = BytesMut::with_capacity(expected_len as usize);
while let Some(chunk) = self.read_next_block().await? {
buf.extend_from_slice(&chunk);
}
Ok(buf.freeze())
}
async fn acquire_worker(&self, addr: &str) -> Result<WorkerClient> {
if let Some(pool) = &self.worker_pool {
pool.acquire(addr).await
} else {
WorkerClient::connect(addr, &self.config).await
}
}
async fn reconnect_worker(
&self,
addr: &str,
stale_generation: Option<u64>,
) -> Result<WorkerClient> {
if let Some(pool) = &self.worker_pool {
match stale_generation {
Some(gen) => pool.reconnect_if_stale(addr, gen).await,
None => pool.reconnect(addr).await,
}
} else {
WorkerClient::connect(addr, &self.config).await
}
}
async fn try_read_block(
&self,
worker: &WorkerClient,
block_id: i64,
plan: &BlockReadPlan,
ufs_options: Option<OpenUfsBlockOptions>,
) -> Result<Bytes> {
let mut block_reader = GrpcBlockReader::open(
worker,
block_id,
plan.offset_in_block as i64,
plan.length as i64,
self.config.chunk_size as i64,
ufs_options,
)
.await?;
block_reader.read_all().await
}
fn resolve_block_id(&self, plan: &BlockReadPlan) -> i64 {
if let Some(fbi) = self
.file_info
.file_block_infos
.get(plan.block_index as usize)
{
if let Some(bi) = &fbi.block_info {
if let Some(id) = bi.block_id {
if id > 0 {
return id;
}
}
}
}
plan.block_id
}
fn build_ufs_read_options(&self, plan: &BlockReadPlan) -> Option<OpenUfsBlockOptions> {
let template = self.ufs_read_options.as_ref()?;
let block_size = template.block_size.unwrap_or(64 * 1024 * 1024);
let offset_in_file = plan.block_index as i64 * block_size;
let mut opts = template.clone();
opts.offset_in_file = Some(offset_in_file);
Some(opts)
}
pub async fn read_file_with_context(ctx: Arc<FileSystemContext>, path: &str) -> Result<Bytes> {
let mut reader = Self::open_with_context(ctx, path).await?;
reader.read_all().await
}
pub async fn read_range_with_context(
ctx: Arc<FileSystemContext>,
path: &str,
offset: u64,
length: u64,
) -> Result<Bytes> {
let mut reader = Self::open_range_with_context(ctx, path, offset, length).await?;
reader.read_all().await
}
pub async fn read_ranges_with_context(
ctx: Arc<FileSystemContext>,
path: &str,
ranges: &[(u64, u64)],
) -> Result<Vec<Bytes>> {
let cfg = ctx.config();
if !cfg.range_coalesce_enabled {
let mut out: Vec<Bytes> = Vec::with_capacity(ranges.len());
for &(off, len) in ranges {
if len == 0 {
out.push(Bytes::new());
continue;
}
let bytes = Self::read_range_with_context(ctx.clone(), path, off, len).await?;
out.push(bytes);
}
return Ok(out);
}
let plan = crate::io::range_coalesce::plan(
ranges,
cfg.range_coalesce_gap_bytes,
cfg.range_coalesce_max_bytes,
);
debug!(
path = %path,
input_count = ranges.len(),
output_count = plan.fetches.len(),
input_bytes = plan.total_input_bytes,
fetch_bytes = plan.total_fetch_bytes,
wasted_bytes = plan.wasted_bytes(),
"range coalesce plan (§B2)"
);
let mut fetch_bufs: Vec<Bytes> = Vec::with_capacity(plan.fetches.len());
for f in &plan.fetches {
let bytes = Self::read_range_with_context(ctx.clone(), path, f.offset, f.len).await?;
if bytes.len() as u64 != f.len {
return Err(Error::BlockIoError {
message: format!(
"read_ranges: merged fetch (offset={}, len={}) returned {} bytes",
f.offset,
f.len,
bytes.len()
),
});
}
fetch_bufs.push(bytes);
}
Ok(Self::splice_from_plan(&plan, &fetch_bufs))
}
fn splice_from_plan(
plan: &crate::io::range_coalesce::CoalescePlan,
fetch_bufs: &[Bytes],
) -> Vec<Bytes> {
debug_assert_eq!(fetch_bufs.len(), plan.fetches.len());
let mut out: Vec<Bytes> = Vec::with_capacity(plan.slices.len());
for s in &plan.slices {
if s.fetch_index == crate::io::range_coalesce::NO_FETCH {
out.push(Bytes::new());
continue;
}
let src = &fetch_bufs[s.fetch_index];
let end = s.offset_in_fetch + s.len;
out.push(src.slice(s.offset_in_fetch..end));
}
out
}
pub fn path(&self) -> &str {
&self.path
}
pub fn file_info(&self) -> &FileInfo {
&self.file_info
}
pub fn file_length(&self) -> u64 {
self.file_info.length.unwrap_or(0) as u64
}
pub fn bytes_read(&self) -> u64 {
self.total_bytes_read
}
pub fn block_count(&self) -> usize {
self.plans.len()
}
pub fn current_block_index(&self) -> usize {
self.current_plan_index
}
pub fn is_complete(&self) -> bool {
self.current_plan_index >= self.plans.len()
}
pub fn offset(&self) -> u64 {
self.offset
}
pub fn length(&self) -> u64 {
self.length
}
}
#[async_trait::async_trait]
impl ExternalRangeReader for GoosefsFileReader {
async fn read_range(&mut self, offset: i64, end: i64) -> Result<Bytes> {
self.read_file_range(offset, end).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::proto::grpc::WorkerNetAddress;
fn make_reader(length: i64, block_size: i64) -> GoosefsFileReader {
let num_blocks = if block_size > 0 {
(length + block_size - 1) / block_size
} else {
0
};
let block_ids: Vec<i64> = (1001..(1001 + num_blocks.max(0))).collect();
let file_info = FileInfo {
length: Some(length),
block_size_bytes: Some(block_size),
block_ids,
completed: Some(true),
folder: Some(false),
ufs_path: Some(String::new()),
..Default::default()
};
let config = GoosefsConfig::new("127.0.0.1:9200");
GoosefsFileReader::build(
&config,
"/synthetic",
Arc::new(file_info),
WorkerRouterView::empty(),
None,
None,
0,
length as u64,
)
.expect("build reader")
}
#[test]
fn test_block_logical_size() {
let bs = 64 * 1024 * 1024i64;
let len = 2 * bs + bs / 2;
let reader = make_reader(len, bs);
assert_eq!(reader.block_logical_size(0), bs);
assert_eq!(reader.block_logical_size(1), bs);
assert_eq!(reader.block_logical_size(2), bs / 2);
assert_eq!(reader.block_logical_size(99), 0);
}
#[test]
fn test_build_defaults_disable_cache_and_sc() {
let reader = make_reader(1024, 1024);
assert!(reader.cache.is_none(), "cache must default to disabled");
assert!(!reader.cache_fill, "fill must default to false");
assert!(
reader.short_circuit.is_none(),
"short-circuit must default to disabled"
);
}
#[test]
fn hr1_non_positive_file_id_is_not_page_cache_eligible() {
use crate::cache::page_cache_eligible;
let reader = make_reader(1024, 1024);
assert_eq!(reader.file_info.file_id, None);
assert!(!page_cache_eligible(reader.file_info.file_id.unwrap_or(0)));
assert!(reader.cache.is_none());
assert!(!page_cache_eligible(0));
assert!(!page_cache_eligible(-7));
assert!(page_cache_eligible(42));
}
#[test]
fn test_worker_addr_formatting() {
let wi = WorkerInfo {
address: Some(WorkerNetAddress {
host: Some("10.0.0.5".to_string()),
rpc_port: Some(9207),
..Default::default()
}),
..Default::default()
};
assert_eq!(
GoosefsFileReader::worker_addr(&wi).unwrap(),
"10.0.0.5:9207"
);
let wi_defaults = WorkerInfo {
address: Some(WorkerNetAddress::default()),
..Default::default()
};
assert_eq!(
GoosefsFileReader::worker_addr(&wi_defaults).unwrap(),
"127.0.0.1:9203"
);
}
#[test]
fn test_worker_addr_missing_address_errors() {
let wi = WorkerInfo {
address: None,
..Default::default()
};
assert!(GoosefsFileReader::worker_addr(&wi).is_err());
}
#[test]
fn splice_from_plan_reconstructs_caller_order_and_bytes() {
use crate::io::range_coalesce::plan;
let inputs = &[(200u64, 10u64), (0, 20), (10, 5), (100, 4)];
let p = plan(inputs, 100, 60);
assert!(!p.fetches.is_empty());
let synth = |off: u64, len: u64| -> Bytes {
let v: Vec<u8> = (0..len).map(|i| ((off + i) & 0xff) as u8).collect();
Bytes::from(v)
};
let fetch_bufs: Vec<Bytes> = p.fetches.iter().map(|f| synth(f.offset, f.len)).collect();
let out = GoosefsFileReader::splice_from_plan(&p, &fetch_bufs);
assert_eq!(out.len(), inputs.len(), "output length must match input");
for (i, &(off, len)) in inputs.iter().enumerate() {
let expected = synth(off, len);
assert_eq!(
out[i].as_ref(),
expected.as_ref(),
"output[{i}] byte mismatch for input ({off},{len})"
);
}
}
#[test]
fn splice_from_plan_preserves_empty_ranges_at_original_position() {
use crate::io::range_coalesce::plan;
let inputs = &[(100u64, 0u64), (0, 10), (7, 0), (50, 5)];
let p = plan(inputs, 0, 4096);
assert_eq!(p.fetches.len(), 2);
let synth = |off: u64, len: u64| -> Bytes {
Bytes::from(
(0..len)
.map(|i| ((off + i) & 0x7f) as u8)
.collect::<Vec<_>>(),
)
};
let fetch_bufs: Vec<Bytes> = p.fetches.iter().map(|f| synth(f.offset, f.len)).collect();
let out = GoosefsFileReader::splice_from_plan(&p, &fetch_bufs);
assert_eq!(out.len(), 4);
assert!(out[0].is_empty(), "empty input must yield empty Bytes");
assert_eq!(out[1].len(), 10);
assert!(out[2].is_empty(), "empty input must yield empty Bytes");
assert_eq!(out[3].len(), 5);
}
#[test]
fn range_coalesce_disabled_by_default() {
let cfg = GoosefsConfig::default();
assert!(
!cfg.range_coalesce_enabled,
"range_coalesce_enabled must default to false (opt-in per §B2)"
);
assert_eq!(cfg.range_coalesce_gap_bytes, 64 * 1024);
assert_eq!(cfg.range_coalesce_max_bytes, 4 * 1024 * 1024);
}
}