use std::io::{Seek, Write};
mod local_file;
mod preallocate;
mod socket;
pub use local_file::LocalFileSink;
pub use socket::{SocketSink, UdpSocketSink};
pub trait SequentialSink: Write + Send {
fn finish(&mut self) -> std::io::Result<()> {
self.flush()
}
}
pub trait RandomAccessSink: SequentialSink + Seek {}
#[allow(dead_code)]
pub(crate) fn open_for_mkv(
dest: &std::path::Path,
size_hint: Option<u64>,
) -> std::io::Result<Box<dyn RandomAccessSink>> {
#[cfg(target_os = "linux")]
{
use crate::platform::fs_type::{FsType, detect};
if detect(dest) == FsType::Nfs {
let wf = match size_hint {
Some(n) => crate::io::WritebackFile::create_with_size_hint(dest, n)?,
None => crate::io::WritebackFile::create(dest)?,
};
return Ok(Box::new(wf));
}
}
#[cfg(not(target_os = "linux"))]
let _ = crate::platform::fs_type::detect;
let sink = match size_hint {
Some(n) => LocalFileSink::with_size_hint(dest, n)?,
None => LocalFileSink::create(dest)?,
};
Ok(Box::new(sink))
}
#[cfg(test)]
mod tests {
use super::*;
fn _assert_is_sequential(_: &mut dyn SequentialSink) {}
fn _assert_is_random_access(_: &mut dyn RandomAccessSink) {}
#[test]
fn concrete_sinks_satisfy_traits() {
let dir = tempfile::tempdir().unwrap();
let mut s = LocalFileSink::create(&dir.path().join("b.bin")).unwrap();
_assert_is_sequential(&mut s);
_assert_is_random_access(&mut s);
let mut wf = crate::io::WritebackFile::create(&dir.path().join("c.bin")).unwrap();
_assert_is_sequential(&mut wf);
_assert_is_random_access(&mut wf);
}
#[test]
fn open_for_mkv_returns_a_random_access_sink() {
let dir = tempfile::tempdir().unwrap();
let p = dir.path().join("c.bin");
let mut sink = open_for_mkv(&p, Some(64 * 1024)).unwrap();
use std::io::{Seek, SeekFrom, Write};
sink.write_all(b"hello").unwrap();
sink.seek(SeekFrom::Start(0)).unwrap();
sink.finish().unwrap();
drop(sink);
let bytes = std::fs::read(&p).unwrap();
assert_eq!(&bytes[..5], b"hello");
}
#[test]
fn finish_through_trait_object_flushes_local_file() {
let dir = tempfile::tempdir().unwrap();
let p = dir.path().join("trait-finish.bin");
let sink = LocalFileSink::create(&p).unwrap();
let mut boxed: Box<dyn SequentialSink> = Box::new(sink);
boxed.write_all(b"buffered-tail").unwrap();
boxed.finish().unwrap();
let bytes = std::fs::read(&p).unwrap();
assert_eq!(&bytes[..], b"buffered-tail");
}
use std::io::{self, Write};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
struct FlushTracker {
flushed: Arc<AtomicBool>,
bytes: Arc<AtomicUsize>,
}
impl Write for FlushTracker {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.bytes.fetch_add(buf.len(), Ordering::SeqCst);
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
self.flushed.store(true, Ordering::SeqCst);
Ok(())
}
}
impl SequentialSink for FlushTracker {}
#[test]
fn default_finish_flushes() {
let flushed = Arc::new(AtomicBool::new(false));
let bytes = Arc::new(AtomicUsize::new(0));
let mut sink = FlushTracker {
flushed: flushed.clone(),
bytes: bytes.clone(),
};
sink.write_all(b"abc").unwrap();
assert!(
!flushed.load(Ordering::SeqCst),
"flush should not run before finish"
);
sink.finish().unwrap();
assert!(
flushed.load(Ordering::SeqCst),
"default SequentialSink::finish must call Write::flush"
);
assert_eq!(bytes.load(Ordering::SeqCst), 3);
}
#[test]
fn open_for_mkv_without_size_hint_is_random_access() {
use std::io::{Seek, SeekFrom};
let dir = tempfile::tempdir().unwrap();
let p = dir.path().join("nohint.bin");
let mut sink = open_for_mkv(&p, None).unwrap();
sink.write_all(b"AAAABBBB").unwrap();
sink.seek(SeekFrom::Start(4)).unwrap();
sink.write_all(b"CCCC").unwrap();
sink.finish().unwrap();
drop(sink);
assert_eq!(std::fs::read(&p).unwrap(), b"AAAACCCC");
}
}