diff --git a/src/common/bitpacker.rs b/src/common/bitpacker.rs index ca9e5e893..a99bc2148 100644 --- a/src/common/bitpacker.rs +++ b/src/common/bitpacker.rs @@ -8,19 +8,17 @@ pub fn compute_num_bits(amplitude: u32) -> u8 { (32u32 - amplitude.leading_zeros()) as u8 } -pub struct BitPacker { - output: TWrite, +pub struct BitPacker { mini_buffer: u64, mini_buffer_written: usize, num_bits: usize, written_size: usize, } -impl BitPacker { +impl BitPacker { - pub fn new(output: TWrite, num_bits: usize) -> BitPacker { + pub fn new(num_bits: usize) -> BitPacker { BitPacker { - output: output, mini_buffer: 0u64, mini_buffer_written: 0, num_bits: num_bits, @@ -28,11 +26,11 @@ impl BitPacker { } } - pub fn write(&mut self, val: u32) -> io::Result<()> { + pub fn write(&mut self, val: u32, output: &mut TWrite) -> io::Result<()> { let val_u64 = val as u64; if self.mini_buffer_written + self.num_bits > 64 { self.mini_buffer |= val_u64.wrapping_shl(self.mini_buffer_written as u32); - self.written_size += self.mini_buffer.serialize(&mut self.output)?; + self.written_size += self.mini_buffer.serialize(output)?; self.mini_buffer = val_u64.wrapping_shr((64 - self.mini_buffer_written) as u32); self.mini_buffer_written = self.mini_buffer_written + (self.num_bits as usize) - 64; } @@ -40,7 +38,7 @@ impl BitPacker { self.mini_buffer |= val_u64 << self.mini_buffer_written; self.mini_buffer_written += self.num_bits; if self.mini_buffer_written == 64 { - self.written_size += self.mini_buffer.serialize(&mut self.output)?; + self.written_size += self.mini_buffer.serialize(output)?; self.mini_buffer_written = 0; self.mini_buffer = 0u64; } @@ -48,37 +46,39 @@ impl BitPacker { Ok(()) } - fn flush(&mut self) -> io::Result<()>{ + fn flush(&mut self, output: &mut TWrite) -> io::Result<()>{ if self.mini_buffer_written > 0 { let num_bytes = (self.mini_buffer_written + 7) / 8; let arr: [u8; 8] = unsafe { mem::transmute::(self.mini_buffer) }; - self.output.write_all(&arr[..num_bytes])?; + output.write_all(&arr[..num_bytes])?; self.written_size += num_bytes; self.mini_buffer_written = 0; } Ok(()) } - pub fn close(mut self) -> io::Result<(TWrite, usize)> { - self.flush()?; - Ok((self.output, self.written_size)) + pub fn close(&mut self, output: &mut TWrite) -> io::Result { + self.flush(output)?; + Ok(self.written_size) } } -pub struct BitUnpacker<'a> { - data: &'a [u8], +pub struct BitUnpacker { num_bits: usize, mask: u32, + data_ptr: *const u8, + data_len: usize, } -impl<'a> BitUnpacker<'a> { - pub fn new(data: &'a [u8], num_bits: usize) -> BitUnpacker<'a> { +impl BitUnpacker { + pub fn new(data: &[u8], num_bits: usize) -> BitUnpacker { BitUnpacker { - data: data, num_bits: num_bits, mask: (1u32 << num_bits) - 1u32, + data_ptr: data.as_ptr(), + data_len: data.len() } } @@ -87,15 +87,17 @@ impl<'a> BitUnpacker<'a> { return 0; } let addr = (idx * self.num_bits) / 8; - let bit_shift = (idx * self.num_bits) - addr * 8; + let bit_shift = idx * self.num_bits - addr * 8; let val_unshifted_unmasked: u64; - if addr + 8 <= self.data.len() { - val_unshifted_unmasked = unsafe { * (self.data.as_ptr().offset(addr as isize) as *const u64) }; + if addr + 8 <= self.data_len { + val_unshifted_unmasked = unsafe { * (self.data_ptr.offset(addr as isize) as *const u64) }; } - else { + else { let mut arr = [0u8; 8]; - for i in 0..self.data.len() - addr { - arr[i] = self.data[addr + i]; + if addr < self.data_len { + for i in 0..self.data_len - addr { + arr[i] = unsafe { *self.data_ptr.offset( (addr + i) as isize) }; + } } val_unshifted_unmasked = unsafe { mem::transmute::<[u8; 8], u64>(arr) }; } @@ -124,7 +126,8 @@ mod test { } fn test_bitpacker_util(len: usize, num_bits: usize) { - let mut bitpacker = BitPacker::new(Vec::new(), num_bits); + let mut data = Vec::new(); + let mut bitpacker = BitPacker::new(num_bits); let max_val: u32 = (1 << num_bits) - 1; let vals: Vec = (0u32..len as u32).map(|i| { if max_val == 0 { @@ -135,9 +138,9 @@ mod test { } }).collect(); for &val in &vals { - bitpacker.write(val).unwrap(); + bitpacker.write(val, &mut data).unwrap(); } - let (data, num_bytes) = bitpacker.close().unwrap(); + let num_bytes = bitpacker.close(&mut data).unwrap(); assert_eq!(num_bytes, (num_bits * len + 7) / 8); assert_eq!(data.len(), num_bytes); let bitunpacker = BitUnpacker::new(&data, num_bits); diff --git a/src/compression/compression.rs b/src/compression/compression.rs index 2c15a73a2..1ab68f98c 100644 --- a/src/compression/compression.rs +++ b/src/compression/compression.rs @@ -71,12 +71,11 @@ mod compression { } let num_bits = compute_num_bits(max_delta); output.write_all(&[num_bits]).unwrap(); - let mut bit_packer = BitPacker::new(output, num_bits as usize); + let mut bit_packer = BitPacker::new(num_bits as usize); for val in &deltas { - bit_packer.write(*val).unwrap(); + bit_packer.write(*val, &mut output).unwrap(); } - let (_, written_size) = bit_packer.close().expect("packing in memory should never fail"); - written_size + 1 + 1 + bit_packer.close(&mut output).expect("packing in memory should never fail") } pub fn uncompress_sorted(compressed_data: &[u8], output: &mut [u32], mut offset: u32) -> usize { @@ -95,12 +94,11 @@ mod compression { let max = vals.iter().cloned().max().expect("compress unsorted called with an empty array"); let num_bits = compute_num_bits(max); output.write_all(&[num_bits]).unwrap(); - let mut bit_packer = BitPacker::new(output, num_bits as usize); + let mut bit_packer = BitPacker::new(num_bits as usize); for val in vals { - bit_packer.write(*val).unwrap(); + bit_packer.write(*val, &mut output).unwrap(); } - let (_, written_size) = bit_packer.close().expect("packing in memory should never fail"); - 1 + written_size + 1 + bit_packer.close(&mut output).expect("packing in memory should never fail") } pub fn uncompress_unsorted(compressed_data: &[u8], output: &mut [u32]) -> usize { diff --git a/src/fastfield/mod.rs b/src/fastfield/mod.rs index eedadae65..b51a5d15b 100644 --- a/src/fastfield/mod.rs +++ b/src/fastfield/mod.rs @@ -18,8 +18,6 @@ pub use self::writer::{U32FastFieldsWriter, U32FastFieldWriter}; pub use self::reader::{U32FastFieldsReader, U32FastFieldReader}; pub use self::serializer::FastFieldSerializer; - - #[cfg(test)] mod tests { use super::*; @@ -79,7 +77,7 @@ mod tests { } let source = directory.open_read(&path).unwrap(); { - assert_eq!(source.len(), 26 as usize); + assert_eq!(source.len(), 20 as usize); } { let fast_field_readers = U32FastFieldsReader::open(source).unwrap(); @@ -112,7 +110,7 @@ mod tests { } let source = directory.open_read(&path).unwrap(); { - assert_eq!(source.len(), 50 as usize); + assert_eq!(source.len(), 45 as usize); } { let fast_field_readers = U32FastFieldsReader::open(source).unwrap(); @@ -147,7 +145,7 @@ mod tests { } let source = directory.open_read(&path).unwrap(); { - assert_eq!(source.len(), 26 as usize); + assert_eq!(source.len(), 18 as usize); } { let fast_field_readers = U32FastFieldsReader::open(source).unwrap(); diff --git a/src/fastfield/reader.rs b/src/fastfield/reader.rs index ee3c8b670..dcf2f2d6a 100644 --- a/src/fastfield/reader.rs +++ b/src/fastfield/reader.rs @@ -12,6 +12,7 @@ use directory::{WritePtr, RAMDirectory, Directory}; use fastfield::FastFieldSerializer; use fastfield::U32FastFieldsWriter; use common::bitpacker::compute_num_bits; +use common::bitpacker::BitUnpacker; lazy_static! { @@ -23,11 +24,9 @@ lazy_static! { pub struct U32FastFieldReader { _data: ReadOnlySource, - data_ptr: *const u8, + bit_unpacker: BitUnpacker, min_val: u32, max_val: u32, - num_bits: u32, - mask: u32, } impl U32FastFieldReader { @@ -47,34 +46,28 @@ impl U32FastFieldReader { pub fn open(data: ReadOnlySource) -> io::Result { let min_val; let amplitude; + let max_val; { let mut cursor = data.as_slice(); min_val = try!(u32::deserialize(&mut cursor)); amplitude = try!(u32::deserialize(&mut cursor)); + max_val = min_val + amplitude; } let num_bits = compute_num_bits(amplitude); - let mask = (1 << num_bits) - 1; - let ptr: *const u8 = &(data.deref()[8 as usize]); + let bit_unpacker = { + let data_arr = &(data.deref()[8..]); + BitUnpacker::new(data_arr, num_bits as usize) + }; Ok(U32FastFieldReader { _data: data, - data_ptr: ptr, + bit_unpacker: bit_unpacker, min_val: min_val, - max_val: min_val + amplitude, - num_bits: num_bits as u32, - mask: mask, + max_val: max_val, }) } pub fn get(&self, doc: DocId) -> u32 { - if self.num_bits == 0u32 { - return self.min_val; - } - let addr = (doc * self.num_bits) / 8; - let bit_shift = (doc * self.num_bits) - addr * 8; //doc - long_addr * self.num_in_pack; - let val_unshifted_unmasked: u64 = unsafe { * (self.data_ptr.offset(addr as isize) as *const u64) }; - let val_shifted = (val_unshifted_unmasked >> bit_shift) as u32; - self.min_val + (val_shifted & self.mask) - + self.min_val + self.bit_unpacker.get(doc as usize) } } diff --git a/src/fastfield/serializer.rs b/src/fastfield/serializer.rs index 06d95462d..a153a1b08 100644 --- a/src/fastfield/serializer.rs +++ b/src/fastfield/serializer.rs @@ -1,9 +1,8 @@ use common::BinarySerializable; use directory::WritePtr; use schema::Field; -use common::bitpacker::compute_num_bits; -use std::io; -use std::io::{Write, Seek, SeekFrom}; +use common::bitpacker::{compute_num_bits, BitPacker}; +use std::io::{self, Write, Seek, SeekFrom}; /// `FastFieldSerializer` is in charge of serializing @@ -30,16 +29,12 @@ pub struct FastFieldSerializer { write: WritePtr, written_size: usize, fields: Vec<(Field, u32)>, - num_bits: u8, min_value: u32, field_open: bool, - - mini_buffer_written: usize, - mini_buffer: u32, + bit_packer: BitPacker, } - impl FastFieldSerializer { /// Constructor pub fn new(mut write: WritePtr) -> io::Result { @@ -49,12 +44,9 @@ impl FastFieldSerializer { write: write, written_size: written_size, fields: Vec::new(), - num_bits: 0u8, min_value: 0, field_open: false, - - mini_buffer_written: 0, - mini_buffer: 0u32, + bit_packer: BitPacker::new(0), }) } @@ -70,37 +62,23 @@ impl FastFieldSerializer { self.written_size += try!(min_value.serialize(write)); let amplitude = max_value - min_value; self.written_size += try!(amplitude.serialize(write)); - self.num_bits = compute_num_bits(amplitude); - if self.num_bits == 0 { - // if num_bits == 0 we make sure that we still write one mini buffer - // so that the reader code does not overflows and does not - // contain a needless if statement. - self.mini_buffer_written += 1; - } + let num_bits = compute_num_bits(amplitude); + self.bit_packer = BitPacker::new(num_bits as usize); + // TODO inspect whether it is required + // if num_bits == 0 { + // // if num_bits == 0 we make sure that we still write one mini buffer + // // so that the reader code does not overflows and does not + // // contain a needless if statement. + // self.mini_buffer_written += 1; + // } Ok(()) } /// Pushes a new value to the currently open u32 fast field. pub fn add_val(&mut self, val: u32) -> io::Result<()> { - let write: &mut Write = &mut self.write; let val_to_write: u32 = val - self.min_value; - if self.mini_buffer_written + self.num_bits as usize > 32 { - self.mini_buffer |= val_to_write.wrapping_shl(self.mini_buffer_written as u32); - self.written_size += try!(self.mini_buffer.serialize(write)); - // overflow of the shift operand is guarded here by the if case. - self.mini_buffer = val_to_write.wrapping_shr(32u32 - self.mini_buffer_written as u32); - self.mini_buffer_written = self.mini_buffer_written + (self.num_bits as usize) - 32 ; - } - else { - self.mini_buffer |= val_to_write << self.mini_buffer_written; - self.mini_buffer_written += self.num_bits as usize; - if self.mini_buffer_written == 32 { - self.written_size += try!(self.mini_buffer.serialize(write)); - self.mini_buffer_written = 0; - self.mini_buffer = 0u32; - } - } + self.bit_packer.write(val_to_write, &mut self.write)?; Ok(()) } @@ -110,15 +88,10 @@ impl FastFieldSerializer { return Err(io::Error::new(io::ErrorKind::Other, "Current field is already closed")); } self.field_open = false; - if self.mini_buffer_written > 0 { - self.mini_buffer_written = 0; - self.written_size += try!(self.mini_buffer.serialize(&mut self.write)); - } // adding some padding to make sure we // can read the last elements with our u64 // cursor - self.written_size += try!(0u32.serialize(&mut self.write)); - self.mini_buffer = 0; + self.written_size += self.bit_packer.close(&mut self.write)?; Ok(()) }