use std::sync::Arc;
use goosefs_sdk::config::GoosefsConfig;
use goosefs_sdk::context::FileSystemContext;
use goosefs_sdk::error::Result;
use goosefs_sdk::io::{GoosefsFileReader, GoosefsFileWriter};
use goosefs_sdk::proto::grpc::file::CreateFilePOptions;
use goosefs_sdk::WritePType;
const TEST_PATH: &str = "/streaming-test/data.bin";
const PAYLOAD_SIZE: usize = 4 * 1024 * 1024;
const BLOCK_SIZE: i64 = 1 * 1024 * 1024;
fn make_payload(size: usize) -> Vec<u8> {
(0..size).map(|i| (i % 251) as u8).collect()
}
fn verify_slice(slice: &[u8], offset: usize) -> bool {
slice
.iter()
.enumerate()
.all(|(i, &b)| b == ((offset + i) % 251) as u8)
}
#[tokio::main]
async fn main() -> Result<()> {
println!("Goosefs Streaming File Read Demo (GoosefsFileReader::read_next_block)");
println!("=====================================================================");
println!("\n0. Creating FileSystemContext...");
let config = GoosefsConfig::new("127.0.0.1:9200");
let ctx: Arc<FileSystemContext> = FileSystemContext::connect(config).await?;
let master = ctx.acquire_master();
println!(" ✅ Context ready");
println!("\n Preparing test directory and file...");
let _ = master.delete(TEST_PATH, false).await;
match master.create_directory("/streaming-test", true).await {
Ok(_) => println!(" Directory /streaming-test created"),
Err(_) => println!(" Directory /streaming-test already exists"),
}
println!(
"\n1. Writing {}-byte test payload ({} MiB)...",
PAYLOAD_SIZE,
PAYLOAD_SIZE / (1024 * 1024)
);
let payload = make_payload(PAYLOAD_SIZE);
let create_opts = CreateFilePOptions {
block_size_bytes: Some(BLOCK_SIZE),
write_type: Some(WritePType::CacheThrough as i32),
recursive: Some(true),
..Default::default()
};
let mut writer =
GoosefsFileWriter::create_with_context(ctx.clone(), TEST_PATH, Some(create_opts)).await?;
for chunk in payload.chunks(256 * 1024) {
writer.write(chunk).await?;
}
writer.close().await?;
println!(
" ✅ Wrote {} bytes (block_size = {} MiB → expected {} blocks)",
writer.bytes_written(),
BLOCK_SIZE / (1024 * 1024),
PAYLOAD_SIZE as i64 / BLOCK_SIZE
);
println!("\n2. FULL read via read_file_with_context (peak memory = O(file))...");
let all = GoosefsFileReader::read_file_with_context(ctx.clone(), TEST_PATH).await?;
println!(
" ✅ read_all returned {} bytes in one contiguous buffer",
all.len()
);
assert_eq!(all.len(), PAYLOAD_SIZE);
assert!(verify_slice(&all, 0), "full-read content mismatch");
drop(all);
println!(
"\n3. STREAMING read via open_with_context + read_next_block \
(peak memory = O(single block))..."
);
let mut reader = GoosefsFileReader::open_with_context(ctx.clone(), TEST_PATH).await?;
println!(
" file_length = {}, block_count = {}",
reader.file_length(),
reader.block_count()
);
let mut total: u64 = 0;
let mut max_block_size: usize = 0;
let mut block_idx = 0usize;
while let Some(chunk) = reader.read_next_block().await? {
let offset = total as usize;
assert!(
verify_slice(&chunk, offset),
"streaming content mismatch at offset {}",
offset
);
total += chunk.len() as u64;
max_block_size = max_block_size.max(chunk.len());
println!(
" block #{:>2}: {:>8} bytes | reader.current_block_index={} reader.bytes_read={}",
block_idx,
chunk.len(),
reader.current_block_index(),
reader.bytes_read()
);
block_idx += 1;
}
println!(
" ✅ Streaming done: {} blocks, {} bytes total, \
largest single block seen = {} bytes",
block_idx, total, max_block_size
);
assert_eq!(total as usize, PAYLOAD_SIZE);
assert_eq!(reader.bytes_read(), PAYLOAD_SIZE as u64);
let range_off: u64 = 1_000_000;
let range_len: u64 = 2_000_000;
println!(
"\n4. STREAMING range read via open_range_with_context \
(offset={}, length={})...",
range_off, range_len
);
let mut range_reader =
GoosefsFileReader::open_range_with_context(ctx.clone(), TEST_PATH, range_off, range_len)
.await?;
let mut range_total: u64 = 0;
let mut range_blocks = 0usize;
while let Some(chunk) = range_reader.read_next_block().await? {
let abs_off = range_off as usize + range_total as usize;
assert!(
verify_slice(&chunk, abs_off),
"range-streaming content mismatch at absolute offset {}",
abs_off
);
range_total += chunk.len() as u64;
range_blocks += 1;
}
println!(
" ✅ Range streaming done: {} blocks, {} bytes (requested {} bytes)",
range_blocks, range_total, range_len
);
assert_eq!(range_total, range_len);
println!("\n5. Post-exhaustion: read_next_block keeps returning None...");
assert!(reader.read_next_block().await?.is_none());
assert!(reader.read_next_block().await?.is_none());
println!(" ✅ Idempotent exhaustion confirmed");
println!("\n6. Cleanup...");
let _ = master.delete(TEST_PATH, false).await;
ctx.close().await?;
println!(" ✅ Context closed");
println!("\n=====================================================================");
println!("✅ Streaming file read demo complete!");
println!("\nTakeaways:");
println!(" • read_file_with_context / read_all → peak memory = O(file)");
println!(" • open_with_context + read_next_block → peak memory = O(single block)");
println!(" • Prefer the streaming API for files that may not fit in RAM,");
println!(" or whenever you process data block-by-block (ETL, parquet, etc.).");
Ok(())
}