use std::fs;
use std::io::{self, Read, Write};
use std::path::Path;
use std::ptr;
use std::time::SystemTime;
use crate::frame::compress::{
lz4f_compress_begin_using_cdict, lz4f_compress_begin_using_dict, LZ4F_VERSION,
};
use crate::frame::header::lz4f_compress_frame_bound;
use crate::frame::types::{
BlockChecksum, BlockMode, BlockSizeId, ContentChecksum, FrameInfo, FrameType, Preferences,
};
use crate::frame::{
lz4f_compress_end, lz4f_compress_frame_using_cdict, lz4f_compress_update,
lz4f_create_compression_context, Lz4FCCtx, Lz4FCDict,
};
use crate::io::file_io::{open_dst_file, open_src_file, NUL_MARK, STDIN_MARK, STDOUT_MARK};
use crate::io::prefs::{display_level, final_time_display, Prefs, KB, LZ4_MAX_DICT_SIZE, MB};
use crate::timefn::get_time;
use crate::util::set_file_stat;
extern "C" {
fn clock() -> libc::clock_t;
}
const CHUNK_SIZE: usize = 4 * MB;
#[derive(Debug, Clone, Copy, Default)]
pub struct CompressStats {
pub bytes_in: u64,
pub bytes_out: u64,
}
pub struct CfcParameters<'a> {
pub prefs: &'a Preferences,
pub cdict: *const Lz4FCDict,
}
unsafe impl<'a> Send for CfcParameters<'a> {}
unsafe impl<'a> Sync for CfcParameters<'a> {}
pub struct CompressResources {
pub src_buffer: Vec<u8>,
pub dst_buffer: Vec<u8>,
pub ctx: Box<Lz4FCCtx>,
pub prepared_prefs: Preferences,
pub cdict: Option<Box<Lz4FCDict>>,
}
unsafe impl Send for CompressResources {}
fn build_preferences(io_prefs: &Prefs) -> Preferences {
let block_size_id = match io_prefs.block_size_id {
4 => BlockSizeId::Max64Kb,
5 => BlockSizeId::Max256Kb,
6 => BlockSizeId::Max1Mb,
_ => BlockSizeId::Max4Mb, };
let block_mode = if io_prefs.block_independence {
BlockMode::Independent
} else {
BlockMode::Linked
};
Preferences {
frame_info: FrameInfo {
block_size_id,
block_mode,
content_checksum_flag: if io_prefs.stream_checksum {
ContentChecksum::Enabled
} else {
ContentChecksum::Disabled
},
block_checksum_flag: if io_prefs.block_checksum {
BlockChecksum::Enabled
} else {
BlockChecksum::Disabled
},
frame_type: FrameType::Frame,
content_size: 0, dict_id: 0,
},
compression_level: 0, auto_flush: true, favor_dec_speed: io_prefs.favor_dec_speed,
}
}
fn effective_block_size(io_prefs: &Prefs) -> usize {
if io_prefs.block_size > 0 {
io_prefs.block_size
} else {
match io_prefs.block_size_id {
4 => 64 * KB,
5 => 256 * KB,
6 => MB,
_ => 4 * MB, }
}
}
fn load_dict_file(dict_filename: &str) -> io::Result<Vec<u8>> {
let circ_size = LZ4_MAX_DICT_SIZE; let mut circular_buf = vec![0u8; circ_size];
let mut dict_end: usize = 0;
let mut dict_len: usize = 0;
let mut reader: Box<dyn Read> = if dict_filename == STDIN_MARK {
Box::new(io::stdin())
} else {
let mut f = fs::File::open(dict_filename).map_err(|e| {
io::Error::new(
e.kind(),
format!("Dictionary error: could not open {}: {}", dict_filename, e),
)
})?;
{
use std::io::Seek;
let _ = f.seek(std::io::SeekFrom::End(-(circ_size as i64)));
}
Box::new(f)
};
loop {
let n = reader.read(&mut circular_buf[dict_end..])?;
if n == 0 {
break; }
dict_end = (dict_end + n) % circ_size;
dict_len += n;
}
if dict_len > LZ4_MAX_DICT_SIZE {
dict_len = LZ4_MAX_DICT_SIZE;
}
let dict_start = (circ_size + dict_end - dict_len) % circ_size;
if dict_start == 0 {
circular_buf.truncate(dict_len);
Ok(circular_buf)
} else {
let first_len = (circ_size - dict_start).min(dict_len);
let second_len = dict_len - first_len;
let mut dict_buf = vec![0u8; dict_len.max(1)];
dict_buf[..first_len].copy_from_slice(&circular_buf[dict_start..dict_start + first_len]);
if second_len > 0 {
dict_buf[first_len..].copy_from_slice(&circular_buf[..second_len]);
}
Ok(dict_buf)
}
}
fn create_cdict(io_prefs: &Prefs) -> io::Result<Option<Box<Lz4FCDict>>> {
if !io_prefs.use_dictionary {
return Ok(None);
}
let dict_filename = io_prefs.dictionary_filename.as_deref().ok_or_else(|| {
io::Error::new(
io::ErrorKind::InvalidInput,
"Dictionary error: no filename provided",
)
})?;
let dict_buf = load_dict_file(dict_filename)?;
let cdict = Lz4FCDict::create(&dict_buf)
.ok_or_else(|| io::Error::other("Dictionary error: could not create CDict"))?;
Ok(Some(cdict))
}
impl CompressResources {
pub fn new(io_prefs: &Prefs) -> io::Result<Self> {
let prepared_prefs = build_preferences(io_prefs);
let ctx = lz4f_create_compression_context(LZ4F_VERSION).map_err(|e| {
io::Error::other(format!(
"Allocation error: can't create LZ4F context: {}",
e
))
})?;
let src_buffer = vec![0u8; CHUNK_SIZE];
let dst_buffer_size = lz4f_compress_frame_bound(CHUNK_SIZE, Some(&prepared_prefs));
let dst_buffer = vec![0u8; dst_buffer_size];
let cdict = create_cdict(io_prefs)?;
Ok(CompressResources {
src_buffer,
dst_buffer,
ctx,
prepared_prefs,
cdict,
})
}
pub fn cdict_ptr(&self) -> *const Lz4FCDict {
self.cdict
.as_deref()
.map_or(ptr::null(), |c| c as *const Lz4FCDict)
}
}
fn read_to_capacity(reader: &mut dyn Read, buf: &mut [u8]) -> io::Result<usize> {
let mut total = 0;
while total < buf.len() {
match reader.read(&mut buf[total..]) {
Ok(0) => break, Ok(n) => total += n,
Err(e) if e.kind() == io::ErrorKind::Interrupted => continue,
Err(e) => return Err(e),
}
}
Ok(total)
}
fn copy_file_stat(src: &str, dst: &str) -> io::Result<()> {
let m = fs::metadata(src)?;
let mtime = m.modified().unwrap_or(SystemTime::UNIX_EPOCH);
#[cfg(unix)]
let (uid, gid, mode) = {
use std::os::unix::fs::MetadataExt;
(m.uid(), m.gid(), m.mode())
};
#[cfg(not(unix))]
let (uid, gid, mode) = (0u32, 0u32, 0o644u32);
set_file_stat(Path::new(dst), mtime, uid, gid, mode)
}
pub fn compress_frame_chunk(
params: &CfcParameters<'_>,
dst: &mut [u8],
src: &[u8],
prefix_data: Option<&[u8]>,
) -> io::Result<usize> {
let mut cctx = lz4f_create_compression_context(LZ4F_VERSION).map_err(|e| {
io::Error::other(format!(
"unable to create a LZ4F compression context: {}",
e
))
})?;
if let Some(prefix) = prefix_data {
lz4f_compress_begin_using_dict(&mut cctx, dst, prefix, Some(params.prefs)).map_err(
|e| {
io::Error::other(format!(
"error initializing LZ4F compression context with prefix: {}",
e
))
},
)?;
} else {
unsafe {
lz4f_compress_begin_using_cdict(&mut cctx, dst, params.cdict, Some(params.prefs))
}
.map_err(|e| {
io::Error::other(format!(
"error initializing LZ4F compression context: {}",
e
))
})?;
}
let c_size = lz4f_compress_update(&mut cctx, dst, src, None).map_err(|e| {
io::Error::other(format!("error compressing with LZ4F_compressUpdate: {}", e))
})?;
Ok(c_size)
}
fn compress_filename_st(
in_stream_size: &mut u64,
ress: &mut CompressResources,
src_filename: &str,
dst_filename: &str,
compression_level: i32,
io_prefs: &Prefs,
) -> io::Result<()> {
let block_size = effective_block_size(io_prefs);
let mut src_reader = open_src_file(src_filename)?;
let mut prefs = ress.prepared_prefs;
prefs.compression_level = compression_level;
if io_prefs.content_size_flag {
let file_size = if src_filename != STDIN_MARK {
fs::metadata(src_filename).map(|m| m.len()).unwrap_or(0)
} else {
0
};
prefs.frame_info.content_size = file_size;
if file_size == 0 {
display_level(3, "Warning : cannot determine input content size \n");
}
}
let dst_file = open_dst_file(dst_filename, io_prefs)?;
let dst_is_stdout = dst_file.is_stdout;
let mut dst_writer: Box<dyn Write> = Box::new(dst_file);
let cdict_ptr = ress.cdict_ptr();
let mut filesize: u64 = 0;
let mut compressedfilesize: u64 = 0;
let mut read_size = read_to_capacity(&mut *src_reader, &mut ress.src_buffer[..block_size])?;
filesize += read_size as u64;
if read_size < block_size {
let c_size = lz4f_compress_frame_using_cdict(
&mut ress.ctx,
&mut ress.dst_buffer,
&ress.src_buffer[..read_size],
cdict_ptr,
Some(&prefs),
)
.map_err(|e| io::Error::other(format!("Compression failed: {}", e)))?;
compressedfilesize = c_size as u64;
display_level(
2,
&format!(
"\rRead : {} MiB ==> {:.2}% ",
filesize >> 20,
compressedfilesize as f64 / (filesize.max(1)) as f64 * 100.0,
),
);
dst_writer
.write_all(&ress.dst_buffer[..c_size])
.map_err(|_| {
io::Error::new(
io::ErrorKind::WriteZero,
"Write error: failed writing single-block compressed frame",
)
})?;
} else {
let header_size = unsafe {
lz4f_compress_begin_using_cdict(
&mut ress.ctx,
&mut ress.dst_buffer,
cdict_ptr,
Some(&prefs),
)
}
.map_err(|e| io::Error::other(format!("File header generation failed: {}", e)))?;
dst_writer
.write_all(&ress.dst_buffer[..header_size])
.map_err(|_| {
io::Error::new(io::ErrorKind::WriteZero, "Write error: cannot write header")
})?;
compressedfilesize += header_size as u64;
while read_size > 0 {
let out_size = lz4f_compress_update(
&mut ress.ctx,
&mut ress.dst_buffer,
&ress.src_buffer[..read_size],
None,
)
.map_err(|e| io::Error::other(format!("Compression failed: {}", e)))?;
compressedfilesize += out_size as u64;
display_level(
2,
&format!(
"\rRead : {} MiB ==> {:.2}% ",
filesize >> 20,
compressedfilesize as f64 / filesize as f64 * 100.0,
),
);
dst_writer
.write_all(&ress.dst_buffer[..out_size])
.map_err(|_| {
io::Error::new(
io::ErrorKind::WriteZero,
"Write error: cannot write compressed block",
)
})?;
read_size = read_to_capacity(&mut *src_reader, &mut ress.src_buffer[..block_size])?;
filesize += read_size as u64;
}
let end_size = lz4f_compress_end(&mut ress.ctx, &mut ress.dst_buffer, None)
.map_err(|e| io::Error::other(format!("End of frame error: {}", e)))?;
dst_writer
.write_all(&ress.dst_buffer[..end_size])
.map_err(|_| {
io::Error::new(
io::ErrorKind::WriteZero,
"Write error: cannot write end of frame",
)
})?;
compressedfilesize += end_size as u64;
}
drop(dst_writer);
if src_filename != STDIN_MARK && !dst_is_stdout && dst_filename != NUL_MARK {
let _ = copy_file_stat(src_filename, dst_filename);
}
if io_prefs.remove_src_file && src_filename != STDIN_MARK {
fs::remove_file(src_filename).map_err(|e| {
io::Error::new(e.kind(), format!("Remove error: {}: {}", src_filename, e))
})?;
}
display_level(2, &format!("\r{:79}\r", ""));
display_level(
2,
&format!(
"Compressed {} bytes into {} bytes ==> {:.2}%\n",
filesize,
compressedfilesize,
compressedfilesize as f64 / filesize.max(1) as f64 * 100.0,
),
);
*in_stream_size = filesize;
Ok(())
}
pub fn compress_filename_ext(
in_stream_size: &mut u64,
ress: &mut CompressResources,
src_filename: &str,
dst_filename: &str,
compression_level: i32,
io_prefs: &Prefs,
) -> io::Result<()> {
compress_filename_st(
in_stream_size,
ress,
src_filename,
dst_filename,
compression_level,
io_prefs,
)
}
pub fn compress_filename(
src: &str,
dst: &str,
compression_level: i32,
prefs: &Prefs,
) -> io::Result<CompressStats> {
let time_start = get_time();
let cpu_start = unsafe { clock() };
let mut ress = CompressResources::new(prefs)?;
let mut processed: u64 = 0;
let result = compress_filename_ext(
&mut processed,
&mut ress,
src,
dst,
compression_level,
prefs,
);
final_time_display(time_start, cpu_start, processed);
result?;
Ok(CompressStats {
bytes_in: processed,
bytes_out: 0,
})
}
pub fn compress_multiple_filenames(
srcs: &[&str],
suffix: &str,
compression_level: i32,
prefs: &Prefs,
) -> io::Result<usize> {
let time_start = get_time();
let cpu_start = unsafe { clock() };
let mut ress = CompressResources::new(prefs)?;
let mut total_processed: u64 = 0;
let mut missed_files: usize = 0;
for &src_name in srcs {
let mut processed: u64 = 0;
let dst_name: String = if suffix == STDOUT_MARK {
STDOUT_MARK.to_owned()
} else {
format!("{}{}", src_name, suffix)
};
if compress_filename_ext(
&mut processed,
&mut ress,
src_name,
&dst_name,
compression_level,
prefs,
)
.is_err()
{
missed_files += 1;
}
total_processed += processed;
}
final_time_display(time_start, cpu_start, total_processed);
Ok(missed_files)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::io::prefs::Prefs;
use tempfile::TempDir;
#[test]
fn compress_resources_new_default_prefs() {
let prefs = Prefs::default();
let ress = CompressResources::new(&prefs).expect("new() should succeed");
assert_eq!(ress.src_buffer.len(), CHUNK_SIZE);
assert!(ress.dst_buffer.len() >= CHUNK_SIZE);
assert!(ress.cdict.is_none());
}
#[test]
fn compress_resources_new_with_dict() {
let dir = TempDir::new().unwrap();
let dict_path = dir.path().join("dict.bin");
std::fs::write(&dict_path, b"hello world dictionary content for testing").unwrap();
let mut prefs = Prefs::default();
prefs.use_dictionary = true;
prefs.dictionary_filename = Some(dict_path.to_str().unwrap().to_owned());
let ress = CompressResources::new(&prefs).expect("new() with dict should succeed");
assert!(ress.cdict.is_some());
}
#[test]
fn load_dict_file_small() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("small.dict");
let data = b"small dictionary data";
std::fs::write(&path, data).unwrap();
let dict = load_dict_file(path.to_str().unwrap()).unwrap();
assert_eq!(dict.as_slice(), data.as_slice());
}
#[test]
fn load_dict_file_large_truncated_to_64kb() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("large.dict");
let data: Vec<u8> = (0u8..=255).cycle().take(96 * 1024).collect();
std::fs::write(&path, &data).unwrap();
let dict = load_dict_file(path.to_str().unwrap()).unwrap();
assert_eq!(dict.len(), LZ4_MAX_DICT_SIZE);
assert_eq!(dict.as_slice(), &data[data.len() - LZ4_MAX_DICT_SIZE..]);
}
#[test]
fn effective_block_size_uses_block_size_when_set() {
let mut p = Prefs::default();
p.block_size = 128 * KB;
assert_eq!(effective_block_size(&p), 128 * KB);
}
#[test]
fn effective_block_size_derives_from_id_when_zero() {
let mut p = Prefs::default();
p.block_size = 0;
p.block_size_id = 4;
assert_eq!(effective_block_size(&p), 64 * KB);
p.block_size_id = 7;
assert_eq!(effective_block_size(&p), 4 * MB);
}
#[test]
fn compress_filename_round_trip_small_file() {
let dir = TempDir::new().unwrap();
let src_path = dir.path().join("input.txt");
let dst_path = dir.path().join("output.lz4");
let original = b"Hello, LZ4 frame format! This is a test of the compression.";
std::fs::write(&src_path, original).unwrap();
let prefs = Prefs::default();
compress_filename(
src_path.to_str().unwrap(),
dst_path.to_str().unwrap(),
1,
&prefs,
)
.expect("compress_filename should succeed");
assert!(dst_path.exists(), "output file must exist");
let compressed = std::fs::read(&dst_path).unwrap();
assert!(compressed.len() >= 7, "must be at least header size");
assert_eq!(
&compressed[..4],
&[0x04, 0x22, 0x4D, 0x18],
"must start with LZ4 magic"
);
let decompressed =
crate::frame::decompress_frame_to_vec(&compressed).expect("decompression must succeed");
assert_eq!(decompressed.as_slice(), original.as_slice());
}
#[test]
fn compress_filename_round_trip_large_file() {
let dir = TempDir::new().unwrap();
let src_path = dir.path().join("large.bin");
let dst_path = dir.path().join("large.lz4");
let original: Vec<u8> = (0u8..=255).cycle().take(200 * 1024).collect();
std::fs::write(&src_path, &original).unwrap();
let mut prefs = Prefs::default();
prefs.block_size_id = 4; prefs.block_size = 64 * KB;
compress_filename(
src_path.to_str().unwrap(),
dst_path.to_str().unwrap(),
1,
&prefs,
)
.expect("compress_filename large should succeed");
let compressed = std::fs::read(&dst_path).unwrap();
let decompressed =
crate::frame::decompress_frame_to_vec(&compressed).expect("decompression must succeed");
assert_eq!(decompressed, original);
}
#[test]
fn compress_multiple_filenames_produces_outputs() {
let dir = TempDir::new().unwrap();
let src1 = dir.path().join("a.txt");
let src2 = dir.path().join("b.txt");
std::fs::write(&src1, b"file a content").unwrap();
std::fs::write(&src2, b"file b content").unwrap();
let prefs = Prefs::default();
let missed = compress_multiple_filenames(
&[src1.to_str().unwrap(), src2.to_str().unwrap()],
".lz4",
1,
&prefs,
)
.expect("compress_multiple_filenames should succeed");
assert_eq!(missed, 0, "no files should be missed");
assert!(dir.path().join("a.txt.lz4").exists());
assert!(dir.path().join("b.txt.lz4").exists());
}
#[test]
fn compress_multiple_filenames_missing_file_counted() {
let prefs = Prefs::default();
let missed = compress_multiple_filenames(
&["/nonexistent/__lz4_missing_file__.txt"],
".lz4",
1,
&prefs,
)
.expect("should return Ok even when some files are missing");
assert_eq!(missed, 1, "one file should be missed");
}
#[test]
fn compress_frame_chunk_returns_nonzero_for_compressible_input() {
let prefs_val = build_preferences(&Prefs::default());
let params = CfcParameters {
prefs: &prefs_val,
cdict: ptr::null(),
};
let src: Vec<u8> = b"abcdefghij".iter().cycle().take(4096).copied().collect();
let mut dst = vec![0u8; lz4f_compress_frame_bound(src.len(), Some(&prefs_val))];
let c_size = compress_frame_chunk(¶ms, &mut dst, &src, None)
.expect("compress_frame_chunk should succeed");
assert!(c_size > 0, "compressed output must be non-empty");
assert!(c_size <= dst.len(), "must not exceed dst capacity");
}
#[test]
fn compress_frame_chunk_with_dict_returns_output() {
let dict_data: Vec<u8> = b"dictionary content"
.iter()
.cycle()
.take(1024)
.copied()
.collect();
let cdict = Lz4FCDict::create(&dict_data).expect("CDict creation failed");
let cdict_ptr: *const Lz4FCDict = &*cdict;
let prefs_val = build_preferences(&Prefs::default());
let params = CfcParameters {
prefs: &prefs_val,
cdict: cdict_ptr,
};
let src: Vec<u8> = b"hello world".iter().cycle().take(512).copied().collect();
let mut dst = vec![0u8; lz4f_compress_frame_bound(src.len(), Some(&prefs_val))];
let c_size = compress_frame_chunk(¶ms, &mut dst, &src, None)
.expect("compress_frame_chunk with dict should succeed");
assert!(c_size > 0);
}
}