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::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::fs::options::ufs_block_length;
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,
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> {
Self::open_inner(ctx, path, None, true).await
}
pub async fn open_range_with_context(
ctx: Arc<FileSystemContext>,
path: &str,
offset: u64,
length: u64,
) -> Result<Self> {
Self::open_inner(ctx, path, Some((offset, length)), true).await
}
async fn open_inner(
ctx: Arc<FileSystemContext>,
path: &str,
range: Option<(u64, u64)>,
use_page_cache: bool,
) -> Result<Self> {
let (file_info, router) = Self::init_with_context(&ctx, path).await?;
let (offset, length) = match range {
Some((offset, length)) => (offset, length),
None => (0, 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()),
offset,
length,
)?;
reader.attach_cache(&ctx, use_page_cache).await;
Ok(reader)
}
async fn init_with_context(
ctx: &Arc<FileSystemContext>,
path: &str,
) -> Result<(Arc<FileInfo>, WorkerRouterView)> {
let mut file_info = ctx.get_file_info_cached(path).await?;
Self::reject_directory(&file_info, path)?;
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)");
let pool = ctx.acquire_worker_pool();
let config = ctx.config();
crate::block::maybe_enrich_file_block_locations(
&mut file_info,
&router,
Some(&pool),
config,
config.check_block_replicas,
)
.await;
Ok((Arc::new(file_info), router))
}
fn reject_directory(file_info: &FileInfo, path: &str) -> Result<()> {
if file_info.folder.unwrap_or(false) {
return Err(Error::OpenDirectory {
path: path.to_string(),
});
}
Ok(())
}
#[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 {
Some(OpenUfsBlockOptions {
ufs_path: file_info.ufs_path.clone(),
offset_in_file: Some(0),
block_size: Some(0),
max_ufs_read_concurrency: Some(
crate::fs::options::InStreamOptions::default().max_ufs_read_concurrency,
),
mount_id: file_info.mount_id,
no_cache: Some(!file_info.cacheable.unwrap_or(true)),
user: None,
caller_type: None,
file_length: file_info.length,
})
}
};
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,
ufs_read_options,
})
}
async fn attach_cache(&mut self, ctx: &Arc<FileSystemContext>, use_page_cache: bool) {
if !use_page_cache {
self.cache = None;
self.cache_fill = false;
return;
}
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, false).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,
positioned: bool,
) -> Result<Bytes> {
let ufs_options = self.build_ufs_read_options(plan);
let locations = self.block_locations(block_id);
let worker_info = self
.router
.select_worker_for_read(
block_id,
locations,
self.config.file_replication_number,
self.config.file_read_max_node_retry,
)
.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_for_read(
block_id,
self.block_locations(block_id),
self.config.file_replication_number,
self.config.file_read_max_node_retry,
)
.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(), positioned)
.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, positioned)
.await
}
Err(e) => Err(e),
}
}
async fn read_file_range(
&self,
abs_offset: i64,
abs_end: i64,
positioned: bool,
) -> 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, positioned).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())
}
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>,
positioned: bool,
) -> Result<Bytes> {
if positioned {
return GrpcBlockReader::positioned_read(
worker,
block_id,
plan.offset_in_block as i64,
plan.length as i64,
self.config.chunk_size as i64,
ufs_options,
)
.await;
}
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 block_locations(&self, block_id: i64) -> &[crate::proto::grpc::BlockLocation] {
self.file_info
.file_block_infos
.iter()
.find_map(|fbi| {
let bi = fbi.block_info.as_ref()?;
if bi.block_id == Some(block_id) {
Some(bi.locations.as_slice())
} else {
None
}
})
.unwrap_or(&[])
}
fn build_ufs_read_options(&self, plan: &BlockReadPlan) -> Option<OpenUfsBlockOptions> {
let template = self.ufs_read_options.as_ref()?;
let nominal_block_size = self.file_info.block_size_bytes.unwrap_or(64 * 1024 * 1024);
let offset_in_file = (plan.block_index as i64).saturating_mul(nominal_block_size);
let mut opts = template.clone();
opts.offset_in_file = Some(offset_in_file);
opts.block_size = Some(match self.file_info.length {
Some(len) => ufs_block_length(len, nominal_block_size, plan.block_index),
None => nominal_block_size,
});
Some(opts)
}
pub async fn read_file_with_context(ctx: Arc<FileSystemContext>, path: &str) -> Result<Bytes> {
let mut reader = Self::open_inner(ctx, path, None, false).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_inner(ctx, path, Some((offset, length)), false).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"
);
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, true).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::fs::options::InStreamOptions;
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 ufs_read_options_carry_max_read_concurrency() {
let bs = 64 * 1024 * 1024i64;
let file_info = FileInfo {
length: Some(bs / 2),
block_size_bytes: Some(bs),
block_ids: vec![1001],
completed: Some(true),
folder: Some(false),
ufs_path: Some("cosn://bucket/blob.bin".to_string()),
mount_id: Some(7),
..Default::default()
};
let reader = GoosefsFileReader::build(
&GoosefsConfig::new("127.0.0.1:9200"),
"/blob.bin",
Arc::new(file_info),
WorkerRouterView::empty(),
None,
None,
0,
(bs / 2) as u64,
)
.expect("build reader");
let opts = reader
.build_ufs_read_options(&reader.plans[0])
.expect("a file with a ufs_path must produce OpenUfsBlockOptions");
assert_eq!(
opts.max_ufs_read_concurrency,
Some(InStreamOptions::default().max_ufs_read_concurrency),
"unset (or 0) makes the Worker refuse every read after the first"
);
assert!(opts.max_ufs_read_concurrency.unwrap() > 0);
}
#[test]
fn ufs_read_options_carry_actual_tail_block_length() {
let bs = 1 << 20i64;
let length = 2 * bs + 100;
let file_info = FileInfo {
length: Some(length),
block_size_bytes: Some(bs),
block_ids: vec![1001, 1002, 1003],
completed: Some(true),
folder: Some(false),
ufs_path: Some("cosn://bucket/tail.lance".to_string()),
mount_id: Some(7),
..Default::default()
};
let reader = GoosefsFileReader::build(
&GoosefsConfig::new("127.0.0.1:9200"),
"/tail.lance",
Arc::new(file_info),
WorkerRouterView::empty(),
None,
None,
0,
length as u64,
)
.expect("build reader");
let seen: Vec<(i64, i64)> = reader
.plans
.iter()
.map(|plan| {
let opts = reader
.build_ufs_read_options(plan)
.expect("a file with a ufs_path must produce OpenUfsBlockOptions");
(opts.offset_in_file.unwrap(), opts.block_size.unwrap())
})
.collect();
assert_eq!(
seen,
vec![(0, bs), (bs, bs), (2 * bs, 100)],
"the tail block must advertise 100 bytes, not the nominal {bs}"
);
}
#[test]
fn ufs_read_options_sub_block_file_reports_file_length() {
let bs = 1 << 20i64;
let file_info = FileInfo {
length: Some(13),
block_size_bytes: Some(bs),
block_ids: vec![1001],
completed: Some(true),
folder: Some(false),
ufs_path: Some("cosn://bucket/latest_version_hint.json".to_string()),
mount_id: Some(7),
..Default::default()
};
let reader = GoosefsFileReader::build(
&GoosefsConfig::new("127.0.0.1:9200"),
"/latest_version_hint.json",
Arc::new(file_info),
WorkerRouterView::empty(),
None,
None,
0,
13,
)
.expect("build reader");
let opts = reader
.build_ufs_read_options(&reader.plans[0])
.expect("a file with a ufs_path must produce OpenUfsBlockOptions");
assert_eq!(opts.block_size, Some(13));
assert_eq!(opts.offset_in_file, Some(0));
}
#[test]
fn test_reject_directory() {
let mut info = FileInfo {
folder: Some(true),
length: Some(0),
..Default::default()
};
let err = GoosefsFileReader::reject_directory(&info, "/d").unwrap_err();
assert!(
matches!(err, Error::OpenDirectory { ref path } if path == "/d"),
"expected OpenDirectory, got {err:?}"
);
info.folder = Some(false);
assert!(GoosefsFileReader::reject_directory(&info, "/f").is_ok());
info.folder = None;
assert!(GoosefsFileReader::reject_directory(&info, "/f").is_ok());
}
#[test]
fn test_build_defaults_disable_cache() {
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");
}
#[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"
);
assert_eq!(cfg.range_coalesce_gap_bytes, 64 * 1024);
assert_eq!(cfg.range_coalesce_max_bytes, 4 * 1024 * 1024);
}
}