use std::io::SeekFrom;
use std::sync::Arc;
use goosefs_sdk::config::GoosefsConfig;
use goosefs_sdk::context::FileSystemContext;
use goosefs_sdk::error::Result;
use goosefs_sdk::fs::options::OpenFileOptions;
use goosefs_sdk::io::{GoosefsAsyncReader, GoosefsFileInStream, GoosefsFileWriter};
use tokio::io::{AsyncReadExt, AsyncSeekExt};
const TEST_PATH: &str = "/async-read-trait/data.bin";
const PAYLOAD_SIZE: usize = 256 * 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(flavor = "current_thread")]
async fn main() -> Result<()> {
let master_addr =
std::env::var("GOOSEFS_MASTER_ADDR").unwrap_or_else(|_| "127.0.0.1:9200".to_string());
let config = GoosefsConfig::new(&master_addr);
println!("== async_read_trait demo ==");
println!("master: {master_addr}");
println!("path: {TEST_PATH}\n");
let ctx: Arc<FileSystemContext> = FileSystemContext::connect(config).await?;
let master = ctx.acquire_master();
let payload = make_payload(PAYLOAD_SIZE);
println!("[1] writing {} bytes via GoosefsFileWriter…", payload.len());
let _ = master.delete(TEST_PATH, false).await;
match master.create_directory("/async-read-trait", true).await {
Ok(_) | Err(_) => {} }
{
let mut writer =
GoosefsFileWriter::create_with_context(ctx.clone(), TEST_PATH, None).await?;
writer.write(&payload).await?;
writer.close().await?;
}
let stream =
GoosefsFileInStream::open_with_context(ctx.clone(), TEST_PATH, OpenFileOptions::default())
.await?;
let mut reader = GoosefsAsyncReader::new(stream);
println!("[2] tokio::io::copy → Vec<u8> (full sequential read via AsyncRead)…");
let mut sink: Vec<u8> = Vec::with_capacity(PAYLOAD_SIZE);
let copied = tokio::io::copy(&mut reader, &mut sink)
.await
.map_err(io_err)?;
assert_eq!(
copied as usize, PAYLOAD_SIZE,
"copy reported {copied} bytes, expected {PAYLOAD_SIZE}"
);
assert_eq!(sink.len(), PAYLOAD_SIZE, "sink size mismatch");
assert!(
verify_slice(&sink, 0),
"byte-for-byte mismatch on full read"
);
println!(" ok — copied {copied} bytes");
println!("[3] seek(Start(60_000)) + read_exact(4096) (random access via AsyncSeek)…");
let pos = reader.seek(SeekFrom::Start(60_000)).await.map_err(io_err)?;
assert_eq!(pos, 60_000, "seek reported wrong position");
let mut chunk = vec![0u8; 4096];
reader.read_exact(&mut chunk).await.map_err(io_err)?;
assert!(
verify_slice(&chunk, 60_000),
"byte-for-byte mismatch on random read"
);
println!(" ok — 4096 bytes at offset 60000 verified");
println!("[4] seek(End(-1024)) + read_to_end (tail read)…");
let tail_off = reader.seek(SeekFrom::End(-1024)).await.map_err(io_err)?;
assert_eq!(
tail_off,
(PAYLOAD_SIZE - 1024) as u64,
"End-relative seek wrong position"
);
let mut tail: Vec<u8> = Vec::new();
let n = reader.read_to_end(&mut tail).await.map_err(io_err)?;
assert_eq!(n, 1024, "tail read returned {n} bytes, expected 1024");
assert!(
verify_slice(&tail, PAYLOAD_SIZE - 1024),
"byte-for-byte mismatch on tail read"
);
println!(" ok — 1024-byte tail verified");
println!("[5] seek(Start(0)) + read_to_end (full round-trip)…");
reader.seek(SeekFrom::Start(0)).await.map_err(io_err)?;
let mut all: Vec<u8> = Vec::with_capacity(PAYLOAD_SIZE);
let n = reader.read_to_end(&mut all).await.map_err(io_err)?;
assert_eq!(n, PAYLOAD_SIZE, "round-trip wrong size");
assert_eq!(all, payload, "round-trip payload mismatch");
println!(" ok — full {PAYLOAD_SIZE}-byte round-trip verified");
let inner = match reader.into_inner() {
Ok(s) => s,
Err(_) => panic!("reader still has an in-flight op after all reads completed"),
};
println!(
"[6] reader.into_inner() ok — underlying stream pos={}",
inner.pos()
);
let _ = master.delete(TEST_PATH, false).await;
println!("\n== all checks passed ==");
Ok(())
}
fn io_err(e: std::io::Error) -> goosefs_sdk::error::Error {
goosefs_sdk::error::Error::Internal {
message: format!("io error: {e}"),
source: None,
}
}