use std::io::SeekFrom;
use std::sync::Arc;
use bytes::{Bytes, BytesMut};
use tracing::{debug, warn};
use crate::block::router::{rpc_endpoint, WorkerRouterView};
use crate::block::short_circuit::{ShortCircuitError, ShortCircuitFactory};
use crate::cache::{page_cache_eligible, CacheManager, ExternalRangeReader};
use crate::client::{WorkerClient, WorkerClientPool, WorkerManagerClient};
use crate::config::GoosefsConfig;
use crate::context::FileSystemContext;
use crate::error::{Error, Result};
use crate::fs::options::InStreamOptions;
use crate::fs::uri_status::URIStatus;
use crate::io::reader::{GrpcBlockReader, ReadTuning};
use crate::proto::proto::dataserver::OpenUfsBlockOptions;
pub const TRANSFER_POSITIONED_READ_THRESHOLD: i64 = 8 * 1024;
#[allow(dead_code)]
const MAX_PREFETCH_WINDOW: i32 = 8;
pub struct GoosefsFileInStream {
status: URIStatus,
config: GoosefsConfig,
options: InStreamOptions,
pos: i64,
file_length: i64,
carry_over: BytesMut,
block_in_stream: Option<GrpcBlockReader>,
block_in_stream_block_id: i64,
#[allow(dead_code)]
cached_positioned_block_id: i64,
router: WorkerRouterView,
worker_pool: Option<Arc<WorkerClientPool>>,
cache: Option<Arc<dyn CacheManager>>,
cache_page_size: u64,
cache_file_id: Arc<str>,
cache_fill: bool,
cache_async_write: bool,
cache_sequential_read: bool,
short_circuit: Option<Arc<ShortCircuitFactory>>,
sc_seq_block: i64,
}
impl GoosefsFileInStream {
pub async fn open(
config: &GoosefsConfig,
path: &str,
options: crate::fs::options::OpenFileOptions,
) -> Result<Self> {
use crate::client::MasterClient;
config
.validate()
.map_err(|e| Error::ConfigError { message: e })?;
let master = MasterClient::connect(config).await?;
let file_info = master.get_status(path).await?;
let status = URIStatus::from_proto(file_info);
if status.is_folder() {
return Err(Error::OpenDirectory {
path: path.to_string(),
});
}
if !status.is_completed() {
return Err(Error::FileIncomplete {
message: format!("{path} is incomplete"),
});
}
let inquire_client = master.inquire_client().clone();
let wm = WorkerManagerClient::connect_with_inquire(config, inquire_client).await?;
let workers = wm.get_worker_info_list().await?;
if workers.is_empty() {
return Err(Error::NoWorkerAvailable {
message: "no workers available for reading".to_string(),
});
}
let router =
WorkerRouterView::from_workers(workers, WorkerRouterView::default_failure_ttl());
let file_length = status.length;
debug!(
path = %path,
file_length = file_length,
block_count = status.block_ids.len(),
"GoosefsFileInStream opened"
);
Ok(Self {
file_length,
status,
config: config.clone(),
options: options.in_stream_options,
pos: 0,
carry_over: BytesMut::new(),
block_in_stream: None,
block_in_stream_block_id: -1,
cached_positioned_block_id: -1,
router,
worker_pool: None, cache: None, cache_page_size: config.client_cache_page_size,
cache_file_id: Arc::from(String::new()),
cache_fill: false,
cache_async_write: false,
cache_sequential_read: false,
short_circuit: None,
sc_seq_block: -1,
})
}
pub async fn open_with_context(
ctx: Arc<FileSystemContext>,
path: &str,
options: crate::fs::options::OpenFileOptions,
) -> Result<Self> {
let config = ctx.config().clone();
config
.validate()
.map_err(|e| Error::ConfigError { message: e })?;
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)");
(*cached).clone()
} 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).clone()
};
let status = URIStatus::from_proto(file_info);
if status.is_folder() {
return Err(Error::OpenDirectory {
path: path.to_string(),
});
}
if !status.is_completed() {
return Err(Error::FileIncomplete {
message: format!("{path} is incomplete"),
});
}
let shared_router = ctx.acquire_router();
let router = WorkerRouterView::from_shared(&shared_router);
let file_length = status.length;
let worker_pool = ctx.acquire_worker_pool();
let short_circuit = ctx.acquire_short_circuit();
let cache = if page_cache_eligible(status.file_id) {
ctx.acquire_cache_manager()
} else {
None
};
let cache_file_id: Arc<str> = Arc::from(status.file_id.to_string());
let cache_fill = cache.is_some()
&& options.in_stream_options.read_type != crate::fs::options::ReadType::NoCache;
if let Some(cache) = &cache {
cache
.on_file_open(
&cache_file_id,
status.length,
status.last_modification_time_ms,
)
.await;
}
debug!(
path = %path,
file_length = file_length,
block_count = status.block_ids.len(),
cache_enabled = cache.is_some(),
"GoosefsFileInStream opened (context mode)"
);
Ok(Self {
file_length,
status,
config: config.clone(),
options: options.in_stream_options,
pos: 0,
carry_over: BytesMut::new(),
block_in_stream: None,
block_in_stream_block_id: -1,
cached_positioned_block_id: -1,
router,
worker_pool: Some(worker_pool),
cache,
cache_page_size: config.client_cache_page_size,
cache_file_id,
cache_fill,
cache_async_write: config.client_cache_async_write_enabled,
cache_sequential_read: config.client_cache_sequential_read_enabled,
short_circuit,
sc_seq_block: -1,
})
}
pub fn pos(&self) -> i64 {
self.pos - self.carry_over.len() as i64
}
pub fn len(&self) -> i64 {
self.file_length
}
pub fn is_empty(&self) -> bool {
self.file_length == 0
}
pub fn is_eof(&self) -> bool {
self.carry_over.is_empty() && self.pos >= self.file_length
}
pub fn remaining(&self) -> i64 {
let raw = (self.file_length - self.pos).max(0);
raw + self.carry_over.len() as i64
}
pub async fn seek(&mut self, pos: i64) -> Result<i64> {
let target = pos.clamp(0, self.file_length);
let user_pos = self.pos();
if target == user_pos {
return Ok(user_pos);
}
self.carry_over.clear();
let seek_dist = (target - self.pos).abs();
let same_block = self.block_index_for_pos(target) == self.block_index_for_pos(self.pos);
if seek_dist < TRANSFER_POSITIONED_READ_THRESHOLD && same_block {
if self.block_in_stream.is_some() {
if target > self.pos {
let skip = (target - self.pos) as usize;
self.skip_bytes(skip).await?;
self.pos = target + self.carry_over.len() as i64;
return Ok(target);
} else {
self.block_in_stream = None;
self.block_in_stream_block_id = -1;
}
}
} else {
self.block_in_stream = None;
self.block_in_stream_block_id = -1;
}
self.pos = target;
Ok(self.pos)
}
pub async fn seek_from(&mut self, seek_from: SeekFrom) -> Result<i64> {
let target = match seek_from {
SeekFrom::Start(n) => n as i64,
SeekFrom::End(n) => self.file_length + n,
SeekFrom::Current(n) => self.pos() + n,
};
self.seek(target).await
}
pub async fn read(&mut self, buf: &mut [u8]) -> Result<usize> {
if buf.is_empty() {
return Ok(0);
}
if !self.carry_over.is_empty() {
let take = self.carry_over.len().min(buf.len());
buf[..take].copy_from_slice(&self.carry_over[..take]);
let _ = self.carry_over.split_to(take);
return Ok(take);
}
if self.pos >= self.file_length {
return Ok(0);
}
if self.cache.is_some() && self.cache_sequential_read {
debug_assert!(
self.carry_over.is_empty(),
"carry_over must be drained before the sequential cache path"
);
let end = (self.pos + buf.len() as i64).min(self.file_length);
let data = self.read_at_cached(self.pos, end).await?;
let n = data.len().min(buf.len());
buf[..n].copy_from_slice(&data[..n]);
self.pos += n as i64;
return Ok(n);
}
let block_idx = self.block_index_for_pos(self.pos);
let block_id = self.block_id_at(block_idx)?;
if self.short_circuit.is_some() {
let use_sc = self.sc_seq_block == block_id
|| (self.block_in_stream_block_id != block_id
&& self.sc_should_use_seq(block_id, block_idx).await);
if use_sc {
if let Some(res) = self.sc_sequential_read(buf, block_idx, block_id).await {
return res;
}
}
}
if self.block_in_stream_block_id != block_id {
let offset_in_block = self.offset_in_block(self.pos);
let remaining_in_block = self.remaining_in_block(self.pos);
let worker = self.connect_worker(block_id).await?;
let worker_generation = worker.generation();
let ufs_opts = self.build_ufs_opts(block_idx);
let tuning = ReadTuning::from_config(&self.config);
let reader_result = GrpcBlockReader::open_sequential(
&worker,
block_id,
offset_in_block,
remaining_in_block,
self.config.chunk_size as i64,
ufs_opts.clone(),
tuning,
)
.await;
let reader = match reader_result {
Ok(r) => r,
Err(e) if e.is_authentication_failed() => {
debug!(
block_id = block_id,
stale_generation = worker_generation,
error = %e,
"auth failed on block reader open, requesting single-flight reconnect"
);
let fresh = self
.reconnect_worker_for_block(block_id, Some(worker_generation))
.await?;
GrpcBlockReader::open_sequential(
&fresh,
block_id,
offset_in_block,
remaining_in_block,
self.config.chunk_size as i64,
ufs_opts,
tuning,
)
.await?
}
Err(e) => return Err(e),
};
self.block_in_stream = Some(reader);
self.block_in_stream_block_id = block_id;
}
let n = self.read_from_sequential_stream(buf).await?;
Ok(n)
}
pub async fn read_at(&mut self, offset: i64, n: usize) -> Result<Bytes> {
if offset >= self.file_length || n == 0 {
return Ok(Bytes::new());
}
let end = (offset + n as i64).min(self.file_length);
if self.cache.is_some() {
return self.read_at_cached(offset, end).await;
}
self.read_external_range(offset, end).await
}
async fn read_external_range(&mut self, offset: i64, end: i64) -> Result<Bytes> {
let first_block_idx = self.block_index_for_pos(offset);
let first_block_end = self.block_start(first_block_idx) + self.status.block_size_bytes;
if end <= first_block_end {
let block_id = self.block_id_at(first_block_idx)?;
let offset_in_block = self.offset_in_block(offset);
let length = end - offset;
if let Some(sc_result) = self
.try_short_circuit_read(block_id, first_block_idx, offset_in_block, length)
.await
{
return sc_result;
}
let data = self
.positioned_read_with_retry(first_block_idx, block_id, offset_in_block, length)
.await?;
if data.is_empty() {
return Err(Error::Internal {
message: format!(
"read_at: positioned_read returned 0 bytes for block {} \
offset_in_block {} length {}",
block_id, offset_in_block, length
),
source: None,
});
}
return Ok(data);
}
let mut result = BytesMut::with_capacity((end - offset) as usize);
let mut cur = offset;
while cur < end {
let block_idx = self.block_index_for_pos(cur);
let block_id = self.block_id_at(block_idx)?;
let offset_in_block = self.offset_in_block(cur);
let block_end = self.block_start(block_idx) + self.status.block_size_bytes;
let read_end = end.min(block_end);
let length = read_end - cur;
let data = match self
.try_short_circuit_read(block_id, block_idx, offset_in_block, length)
.await
{
Some(Ok(b)) => b,
Some(Err(e)) => return Err(e),
None => {
self.positioned_read_with_retry(block_idx, block_id, offset_in_block, length)
.await?
}
};
let advanced = data.len() as i64;
result.extend_from_slice(&data);
if advanced == 0 {
return Err(Error::Internal {
message: format!(
"read_at: positioned_read returned 0 bytes for block {} \
offset_in_block {} length {} (cur={}, end={})",
block_id, offset_in_block, length, cur, end
),
source: None,
});
}
cur += advanced;
}
Ok(result.freeze())
}
async fn read_at_cached(&mut self, offset: i64, end: i64) -> Result<Bytes> {
let cache = match self.cache.clone() {
Some(c) => c,
None => return self.read_external_range(offset, end).await,
};
let page_size = self.cache_page_size;
let file_id = self.cache_file_id.clone();
let file_length = self.file_length;
let fill_mode = self.cache_fill_mode();
crate::cache::read_through_cache(
&cache,
self,
&file_id,
page_size,
file_length,
offset,
end,
fill_mode,
)
.await
}
fn cache_fill_mode(&self) -> crate::cache::FillMode {
if !self.cache_fill {
crate::cache::FillMode::None
} else if self.cache_async_write {
crate::cache::FillMode::Async
} else {
crate::cache::FillMode::Sync
}
}
async fn positioned_read_with_retry(
&mut self,
block_idx: usize,
block_id: i64,
offset_in_block: i64,
length: i64,
) -> Result<Bytes> {
let worker = self.connect_worker(block_id).await?;
let worker_generation = worker.generation();
let ufs_opts = self.build_ufs_opts(block_idx);
let read_result = GrpcBlockReader::positioned_read(
&worker,
block_id,
offset_in_block,
length,
self.config.chunk_size as i64,
ufs_opts.clone(),
)
.await;
match read_result {
Ok(d) => Ok(d),
Err(e) if e.is_authentication_failed() => {
debug!(
block_id = block_id,
stale_generation = worker_generation,
error = %e,
"auth failed on positioned read, requesting single-flight reconnect"
);
let fresh = self
.reconnect_worker_for_block(block_id, Some(worker_generation))
.await?;
GrpcBlockReader::positioned_read(
&fresh,
block_id,
offset_in_block,
length,
self.config.chunk_size as i64,
ufs_opts,
)
.await
}
Err(e) => Err(e),
}
}
fn block_logical_size(&self, block_idx: usize) -> i64 {
let bs = self.status.block_size_bytes;
if bs <= 0 {
return self.file_length.max(0);
}
let start = self.block_start(block_idx);
(self.file_length - start).clamp(0, bs)
}
async fn try_short_circuit_read(
&self,
block_id: i64,
block_idx: usize,
offset_in_block: i64,
length: i64,
) -> 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(block_idx);
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(offset_in_block as usize, 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"
);
factory.invalidate(block_id).await;
None
}
}
}
async fn sc_should_use_seq(&self, block_id: i64, block_idx: usize) -> bool {
match &self.short_circuit {
Some(f) => {
f.should_use(block_id, self.block_logical_size(block_idx))
.await
}
None => false,
}
}
async fn sc_sequential_read(
&mut self,
buf: &mut [u8],
block_idx: usize,
block_id: i64,
) -> Option<Result<usize>> {
let factory = self.short_circuit.clone()?;
let block_size = self.block_logical_size(block_idx);
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 sequential open failed, falling back to gRPC"
);
self.sc_seq_block = -1;
return None;
}
};
let off = self.offset_in_block(self.pos) as usize;
let remaining = self.remaining_in_block(self.pos).max(0) as usize;
let n = remaining.min(buf.len());
if n == 0 {
self.sc_seq_block = -1;
return Some(Ok(0));
}
match reader.read(off, n) {
Ok(slice) => {
buf[..n].copy_from_slice(slice);
self.block_in_stream = None;
self.block_in_stream_block_id = -1;
self.sc_seq_block = block_id;
self.pos += n as i64;
Some(Ok(n))
}
Err(ShortCircuitError::OutOfRange {
off,
len,
file_size,
}) => Some(Err(Error::InvalidArgument {
message: format!(
"short-circuit sequential 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 sequential read failed, falling back to gRPC"
);
factory.invalidate(block_id).await;
self.sc_seq_block = -1;
None
}
}
}
pub async fn read_all(&mut self) -> Result<Bytes> {
let remaining = self.remaining() as usize;
let mut buf = BytesMut::with_capacity(remaining);
let scratch_cap = (self.config.chunk_size as usize).clamp(8 * 1024, 64 * 1024);
let mut tmp = vec![0u8; scratch_cap];
loop {
let n = self.read(&mut tmp).await?;
if n == 0 {
break;
}
buf.extend_from_slice(&tmp[..n]);
}
Ok(buf.freeze())
}
fn block_index_for_pos(&self, offset: i64) -> usize {
if self.status.block_size_bytes <= 0 {
return 0;
}
(offset / self.status.block_size_bytes) as usize
}
fn block_start(&self, idx: usize) -> i64 {
idx as i64 * self.status.block_size_bytes
}
fn offset_in_block(&self, offset: i64) -> i64 {
if self.status.block_size_bytes <= 0 {
return offset;
}
offset % self.status.block_size_bytes
}
fn remaining_in_block(&self, offset: i64) -> i64 {
let block_idx = self.block_index_for_pos(offset);
let block_end = self.block_start(block_idx) + self.status.block_size_bytes;
block_end.min(self.file_length) - offset
}
fn block_id_at(&self, block_idx: usize) -> Result<i64> {
if let Some(&id) = self.status.block_ids.get(block_idx) {
if id > 0 {
return Ok(id);
}
}
let id = self
.status
.block_infos()
.values()
.find(|fbi| {
fbi.offset
.is_some_and(|off| off == self.block_start(block_idx))
})
.and_then(|fbi| fbi.block_info.as_ref())
.and_then(|bi| bi.block_id)
.unwrap_or(-1);
if id > 0 {
Ok(id)
} else {
Err(Error::Internal {
message: format!("no valid block_id at index {block_idx}"),
source: None,
})
}
}
async fn connect_worker(&mut self, block_id: i64) -> Result<WorkerClient> {
let worker_info = self.router.select_worker(block_id).await?;
let addr = worker_info
.address
.as_ref()
.ok_or_else(|| Error::Internal {
message: "worker has no address".to_string(),
source: None,
})?;
let worker_addr = rpc_endpoint(addr);
let result = if let Some(pool) = &self.worker_pool {
pool.acquire(&worker_addr).await
} else {
WorkerClient::connect(&worker_addr, &self.config).await
};
match result {
Ok(w) => Ok(w),
Err(e) => {
if matches!(e, Error::AuthenticationFailed { .. }) {
debug!(
worker = %worker_addr,
error = %e,
"authentication failed on acquire, reconnecting with fresh credentials"
);
if let Some(pool) = &self.worker_pool {
return pool.reconnect(&worker_addr).await;
}
return WorkerClient::connect(&worker_addr, &self.config).await;
}
self.router.mark_failed(addr);
if let Some(pool) = &self.worker_pool {
pool.invalidate(&worker_addr).await;
}
warn!(worker = %worker_addr, error = %e, "worker connect failed, retrying");
let retry_info = self.router.select_worker(block_id).await?;
let retry_addr_info =
retry_info.address.as_ref().ok_or_else(|| Error::Internal {
message: "retry worker has no address".to_string(),
source: None,
})?;
let retry_addr = rpc_endpoint(retry_addr_info);
if let Some(pool) = &self.worker_pool {
pool.acquire(&retry_addr).await
} else {
WorkerClient::connect(&retry_addr, &self.config).await
}
}
}
}
async fn reconnect_worker_for_block(
&mut self,
block_id: i64,
stale_generation: Option<u64>,
) -> Result<WorkerClient> {
let worker_info = self.router.select_worker(block_id).await?;
let addr = worker_info
.address
.as_ref()
.ok_or_else(|| Error::Internal {
message: "worker has no address".to_string(),
source: None,
})?;
let worker_addr = rpc_endpoint(addr);
if let Some(pool) = &self.worker_pool {
match stale_generation {
Some(gen) => pool.reconnect_if_stale(&worker_addr, gen).await,
None => pool.reconnect(&worker_addr).await,
}
} else {
WorkerClient::connect(&worker_addr, &self.config).await
}
}
fn build_ufs_opts(&self, block_idx: usize) -> Option<OpenUfsBlockOptions> {
let ufs_path = self.status.ufs_path.as_str();
if ufs_path.is_empty() {
return None;
}
let block_size = self.status.block_size_bytes;
let offset_in_file = block_idx as i64 * block_size;
Some(OpenUfsBlockOptions {
ufs_path: Some(ufs_path.to_string()),
offset_in_file: Some(offset_in_file),
block_size: Some(block_size),
max_ufs_read_concurrency: Some(self.options.max_ufs_read_concurrency),
mount_id: Some(self.status.mount_id),
no_cache: Some(!self.status.cacheable),
user: None,
caller_type: None,
})
}
async fn skip_bytes(&mut self, mut skip: usize) -> Result<()> {
let stream = match self.block_in_stream.as_mut() {
Some(s) => s,
None => return Ok(()),
};
while skip > 0 {
match stream.read_chunk().await? {
Some(data) => {
if data.len() > skip {
self.carry_over.extend_from_slice(&data[skip..]);
skip = 0;
} else {
skip -= data.len();
}
}
None => break,
}
}
Ok(())
}
async fn read_from_sequential_stream(&mut self, buf: &mut [u8]) -> Result<usize> {
debug_assert!(
self.carry_over.is_empty(),
"carry_over must be drained before pulling a fresh chunk"
);
let stream = match self.block_in_stream.as_mut() {
Some(s) => s,
None => return Ok(0),
};
match stream.read_chunk().await? {
Some(data) => {
let n = data.len().min(buf.len());
buf[..n].copy_from_slice(&data[..n]);
if data.len() > n {
self.carry_over.extend_from_slice(&data[n..]);
debug!(
chunk_len = data.len(),
copied = n,
carried = self.carry_over.len(),
"chunk larger than caller buffer — overflow parked in carry_over"
);
}
self.pos += data.len() as i64;
if stream.is_complete() {
self.block_in_stream = None;
self.block_in_stream_block_id = -1;
}
Ok(n)
}
None => {
self.block_in_stream = None;
self.block_in_stream_block_id = -1;
Ok(0)
}
}
}
}
#[async_trait::async_trait]
impl ExternalRangeReader for GoosefsFileInStream {
async fn read_range(&mut self, offset: i64, end: i64) -> Result<Bytes> {
self.read_external_range(offset, end).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::fs::uri_status::URIStatus;
use crate::proto::grpc::file::FileInfo;
fn make_status(length: i64, block_size: i64) -> URIStatus {
let num_blocks = (length + block_size - 1) / block_size;
let block_ids: Vec<i64> = (1001..(1001 + num_blocks)).collect();
let fi = FileInfo {
length: Some(length),
block_size_bytes: Some(block_size),
block_ids: block_ids.clone(),
completed: Some(true),
folder: Some(false),
ufs_path: Some(String::new()), ..Default::default()
};
URIStatus::from_proto(fi)
}
fn make_stream(status: URIStatus) -> GoosefsFileInStream {
let config = crate::config::GoosefsConfig::new("127.0.0.1:9200");
let file_length = status.length;
GoosefsFileInStream {
file_length,
status,
config,
options: InStreamOptions::default(),
pos: 0,
carry_over: BytesMut::new(),
block_in_stream: None,
block_in_stream_block_id: -1,
cached_positioned_block_id: -1,
router: WorkerRouterView::empty(),
worker_pool: None,
cache: None,
cache_page_size: 1024 * 1024,
cache_file_id: Arc::from(String::new()),
cache_fill: false,
cache_async_write: false,
cache_sequential_read: false,
short_circuit: None,
sc_seq_block: -1,
}
}
#[test]
fn test_block_index_calculation() {
let bs = 64 * 1024 * 1024i64;
let len = 3 * bs;
let status = make_status(len, bs);
let stream = make_stream(status);
assert_eq!(stream.block_index_for_pos(0), 0);
assert_eq!(stream.block_index_for_pos(bs - 1), 0);
assert_eq!(stream.block_index_for_pos(bs), 1);
assert_eq!(stream.block_index_for_pos(2 * bs), 2);
}
#[test]
fn hr1_non_positive_file_id_disables_page_cache_eligibility() {
use crate::cache::page_cache_eligible;
let mut status = make_status(1024, 1024);
status.file_id = 0;
assert!(!page_cache_eligible(status.file_id));
let stream = make_stream(status);
assert!(
stream.cache.is_none(),
"synthetic stream without open must not enable cache"
);
let mut status_neg = make_status(1024, 1024);
status_neg.file_id = -1;
assert!(!page_cache_eligible(status_neg.file_id));
let mut status_ok = make_status(1024, 1024);
status_ok.file_id = 1001;
assert!(page_cache_eligible(status_ok.file_id));
}
#[test]
fn test_offset_in_block() {
let bs = 64 * 1024 * 1024i64;
let status = make_status(2 * bs, bs);
let stream = make_stream(status);
assert_eq!(stream.offset_in_block(0), 0);
assert_eq!(stream.offset_in_block(100), 100);
assert_eq!(stream.offset_in_block(bs), 0);
assert_eq!(stream.offset_in_block(bs + 42), 42);
}
#[test]
fn test_remaining_in_block() {
let bs = 64 * 1024 * 1024i64;
let status = make_status(2 * bs, bs);
let stream = make_stream(status);
assert_eq!(stream.remaining_in_block(0), bs);
assert_eq!(stream.remaining_in_block(bs - 100), 100);
assert_eq!(stream.remaining_in_block(bs), bs);
}
#[test]
fn test_block_id_at() {
let bs = 64 * 1024 * 1024i64;
let status = make_status(2 * bs, bs);
let stream = make_stream(status);
assert_eq!(stream.block_id_at(0).unwrap(), 1001);
assert_eq!(stream.block_id_at(1).unwrap(), 1002);
assert!(stream.block_id_at(99).is_err()); }
#[test]
fn test_is_eof() {
let bs = 1024i64;
let status = make_status(bs, bs);
let mut stream = make_stream(status);
assert!(!stream.is_eof());
stream.pos = bs;
assert!(stream.is_eof());
}
#[test]
fn test_remaining() {
let bs = 1024i64;
let status = make_status(bs, bs);
let mut stream = make_stream(status);
assert_eq!(stream.remaining(), bs);
stream.pos = 100;
assert_eq!(stream.remaining(), bs - 100);
stream.pos = bs;
assert_eq!(stream.remaining(), 0);
}
#[test]
fn test_pos_accounts_for_carry_over() {
let bs = 1024i64;
let status = make_status(bs, bs);
let mut stream = make_stream(status);
stream.pos = 200;
stream.carry_over.extend_from_slice(&[0u8; 50]);
assert_eq!(stream.pos(), 150);
assert_eq!(stream.remaining(), bs - 150);
assert!(!stream.is_eof());
stream.carry_over.clear();
assert_eq!(stream.pos(), 200);
assert_eq!(stream.remaining(), bs - 200);
}
#[test]
fn test_is_eof_with_carry_over() {
let bs = 1024i64;
let status = make_status(bs, bs);
let mut stream = make_stream(status);
stream.pos = bs;
stream.carry_over.extend_from_slice(&[7u8; 20]);
assert!(!stream.is_eof(), "carry_over still has bytes — not EOF");
stream.carry_over.clear();
assert!(
stream.is_eof(),
"chunk reader done and carry_over drained — EOF"
);
}
#[test]
fn test_legacy_mode_no_pool() {
let bs = 1024i64;
let status = make_status(bs, bs);
let stream = make_stream(status);
assert!(
stream.worker_pool.is_none(),
"legacy mode should have no pool"
);
assert!(
stream.short_circuit.is_none(),
"legacy mode should not enable short-circuit"
);
}
#[test]
fn test_block_logical_size() {
let bs = 64 * 1024 * 1024i64;
let len = 2 * bs + bs / 2;
let status = make_status(len, bs);
let stream = make_stream(status);
assert_eq!(stream.block_logical_size(0), bs);
assert_eq!(stream.block_logical_size(1), bs);
assert_eq!(stream.block_logical_size(2), bs / 2);
assert_eq!(stream.block_logical_size(99), 0);
}
}