use std::io;
use std::mem;
use std::sync::Arc;
use bytes::Buf;
use cheetah_string::CheetahString;
use rocketmq_common::common::hasher::string_hasher::JavaStringHasher;
use tracing::info;
use tracing::warn;
use crate::index::index_header::IndexHeader;
use crate::index::index_header::INDEX_HEADER_SIZE;
use crate::log_file::mapped_file::default_mapped_file_impl::DefaultMappedFile;
use crate::log_file::mapped_file::MappedFile;
const HASH_SLOT_SIZE: usize = 4;
const INDEX_SIZE: usize = 20;
const INVALID_INDEX: i32 = 0;
pub const DEFAULT_HASH_SLOT_NUM: usize = 5_000_000;
pub const DEFAULT_INDEX_NUM: usize = 20_000_000;
pub struct IndexFile {
hash_slot_num: usize,
index_num: usize,
file_total_size: usize,
mapped_file: Arc<DefaultMappedFile>,
index_header: IndexHeader,
}
impl PartialEq for IndexFile {
fn eq(&self, other: &Self) -> bool {
std::ptr::eq(self as *const IndexFile, other as *const IndexFile)
}
}
impl IndexFile {
pub fn new(
file_name: &str,
hash_slot_num: usize,
index_num: usize,
end_phy_offset: i64,
end_timestamp: i64,
) -> IndexFile {
Self::try_new(file_name, hash_slot_num, index_num, end_phy_offset, end_timestamp)
.expect("Create index file failed")
}
pub fn try_new(
file_name: &str,
hash_slot_num: usize,
index_num: usize,
end_phy_offset: i64,
end_timestamp: i64,
) -> io::Result<IndexFile> {
let file_total_size = index_file_total_size(hash_slot_num, index_num)?;
let mapped_file = Arc::new(DefaultMappedFile::try_new(
CheetahString::from_slice(file_name),
file_total_size as u64,
)?);
let index_header = IndexHeader::new(mapped_file.clone());
let index_file = IndexFile {
hash_slot_num,
index_num,
file_total_size,
mapped_file,
index_header,
};
if end_phy_offset > 0 {
index_file.index_header.set_begin_phy_offset(end_phy_offset);
index_file.index_header.set_end_phy_offset(end_phy_offset);
}
if end_timestamp > 0 {
index_file.index_header.set_begin_timestamp(end_timestamp);
index_file.index_header.set_end_timestamp(end_timestamp);
}
Ok(index_file)
}
#[inline]
pub fn get_file_name(&self) -> &CheetahString {
self.mapped_file.get_file_name()
}
#[inline]
pub fn get_file_size(&self) -> usize {
self.file_total_size
}
#[inline]
pub fn load(&self) {
self.index_header.load();
}
#[inline]
pub fn shutdown(&self) {
self.flush();
}
pub fn flush(&self) {
let begin_time = std::time::Instant::now();
if self.mapped_file.hold() {
self.index_header.update_byte_buffer();
self.mapped_file.flush(0);
self.mapped_file.release();
info!("flush index file elapsed time(ms) {}", begin_time.elapsed().as_millis());
}
}
#[inline]
pub fn is_write_full(&self) -> bool {
self.index_header.get_index_count() >= self.index_num as i32
}
#[inline]
pub fn destroy(&self, interval_forcibly: u64) -> bool {
self.mapped_file.destroy(interval_forcibly)
}
pub fn put_key(&self, key: &str, phy_offset: i64, store_timestamp: i64) -> bool {
if self.index_header.get_index_count() < self.index_num as i32 {
let hash_code = self.index_key_hash_method(key);
let slot_pos = hash_code as usize % self.hash_slot_num;
let abs_slot_pos = INDEX_HEADER_SIZE + slot_pos * HASH_SLOT_SIZE;
let mapped_file = self.mapped_file.get_mapped_file_mut();
let mut slot_value = mapped_file.get(abs_slot_pos..abs_slot_pos + 4).unwrap().get_i32();
if slot_value <= INVALID_INDEX || slot_value > self.index_header.get_index_count() {
slot_value = INVALID_INDEX;
}
let mut time_diff = store_timestamp - self.index_header.get_begin_timestamp();
time_diff /= 1000;
if self.index_header.get_begin_timestamp() <= 0 {
time_diff = 0;
} else if time_diff > i32::MAX as i64 {
time_diff = i32::MAX as i64;
} else if time_diff < 0 {
time_diff = 0;
}
let abs_index_pos = INDEX_HEADER_SIZE
+ self.hash_slot_num * HASH_SLOT_SIZE
+ self.index_header.get_index_count() as usize * INDEX_SIZE;
self.mapped_file
.write_bytes_segment(&hash_code.to_be_bytes(), abs_index_pos, 0, mem::size_of::<i32>());
self.mapped_file.write_bytes_segment(
&phy_offset.to_be_bytes(),
abs_index_pos + 4,
0,
mem::size_of::<i64>(),
);
self.mapped_file.write_bytes_segment(
&(time_diff as i32).to_be_bytes(),
abs_index_pos + 4 + 8,
0,
mem::size_of::<i32>(),
);
self.mapped_file.write_bytes_segment(
&slot_value.to_be_bytes(),
abs_index_pos + 4 + 8 + 4,
0,
mem::size_of::<i32>(),
);
self.mapped_file.write_bytes_segment(
&self.index_header.get_index_count().to_be_bytes(),
abs_slot_pos,
0,
mem::size_of::<i32>(),
);
if self.index_header.get_index_count() <= 1 {
self.index_header.set_begin_phy_offset(phy_offset);
self.index_header.set_begin_timestamp(store_timestamp);
}
if slot_value == INVALID_INDEX {
self.index_header.inc_hash_slot_count();
}
self.index_header.inc_index_count();
self.index_header.set_end_phy_offset(phy_offset);
self.index_header.set_end_timestamp(store_timestamp);
true
} else {
warn!(
"Over index file capacity: index count = {}; index max num = {}",
self.index_header.get_index_count(),
self.index_num
);
false
}
}
pub fn index_key_hash_method(&self, key: &str) -> i32 {
let hash_code = JavaStringHasher::hash_str(key);
if hash_code == i32::MIN {
0
} else {
hash_code.abs()
}
}
#[inline]
pub fn get_begin_timestamp(&self) -> i64 {
self.index_header.get_begin_timestamp()
}
#[inline]
pub fn get_end_timestamp(&self) -> i64 {
self.index_header.get_end_timestamp()
}
#[inline]
pub fn get_end_phy_offset(&self) -> i64 {
self.index_header.get_end_phy_offset()
}
#[inline]
pub fn has_entries(&self) -> bool {
self.index_header.get_index_count() > 1
}
pub fn is_time_matched(&self, begin: i64, end: i64) -> bool {
let begin_timestamp = self.index_header.get_begin_timestamp();
let end_timestamp = self.index_header.get_end_timestamp();
begin < begin_timestamp && end > end_timestamp
|| begin >= begin_timestamp && begin <= end_timestamp
|| end >= begin_timestamp && end <= end_timestamp
}
pub fn select_phy_offset(&self, phy_offsets: &mut Vec<i64>, key: &str, max_num: usize, begin: i64, end: i64) {
if !self.mapped_file.hold() {
return;
}
(|| {
let key_hash = self.index_key_hash_method(key);
let slot_pos = key_hash as usize % self.hash_slot_num;
let abs_slot_pos = INDEX_HEADER_SIZE + slot_pos * HASH_SLOT_SIZE;
let mut buffer = match self.mapped_file.get_slice(abs_slot_pos, abs_slot_pos + 4) {
None => return,
Some(value) => value,
};
let slot_value = buffer.get_i32();
if slot_value <= INVALID_INDEX
|| slot_value > self.index_header.get_index_count()
|| self.index_header.get_index_count() <= 1
{
return;
}
let mut next_index_to_read = slot_value;
while phy_offsets.len() < max_num {
let abs_index_pos =
INDEX_HEADER_SIZE + self.hash_slot_num * HASH_SLOT_SIZE + next_index_to_read as usize * INDEX_SIZE;
let buffer = match self.mapped_file.get_slice(abs_index_pos, abs_index_pos + INDEX_SIZE) {
None => break,
Some(buf) => buf,
};
let key_hash_read = (&buffer[0..4]).get_i32();
let phy_offset_read = (&buffer[4..12]).get_i64();
let time_diff = (&buffer[12..16]).get_i32();
let prev_index_read = (&buffer[16..20]).get_i32();
if time_diff < 0 {
break;
}
let time_read = self.index_header.get_begin_timestamp() + time_diff as i64 * 1000;
if key_hash == key_hash_read && (time_read >= begin && time_read <= end) {
phy_offsets.push(phy_offset_read);
}
if prev_index_read <= INVALID_INDEX
|| prev_index_read > self.index_header.get_index_count()
|| prev_index_read == next_index_to_read
|| time_read < begin
{
break;
}
next_index_to_read = prev_index_read;
}
})();
self.mapped_file.release();
}
}
fn index_file_total_size(hash_slot_num: usize, index_num: usize) -> io::Result<usize> {
if hash_slot_num == 0 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"index hash slot number must be positive",
));
}
if index_num == 0 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"index entry number must be positive",
));
}
let hash_slots_size = hash_slot_num
.checked_mul(HASH_SLOT_SIZE)
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "index hash slot section size overflow"))?;
let indexes_size = index_num
.checked_mul(INDEX_SIZE)
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "index entry section size overflow"))?;
INDEX_HEADER_SIZE
.checked_add(hash_slots_size)
.and_then(|size| size.checked_add(indexes_size))
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "index file total size overflow"))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_index_key_hash_method_consistency() {
let file = create_test_index_file("20000000000000");
assert_eq!(file.index_key_hash_method("hello"), 99162322);
assert_eq!(file.index_key_hash_method(""), 0);
assert_eq!(file.index_key_hash_method("test"), 3556498);
let hash_result = file.index_key_hash_method("some_key_that_produces_min");
assert!(hash_result >= 0, "Hash should be positive after abs()");
}
#[test]
fn test_put_key_basic() {
let file = create_test_index_file("20000000000001");
assert!(file.put_key("key1", 1000, 1000000000000));
assert_eq!(file.index_header.get_index_count(), 2);
assert!(file.put_key("key2", 2000, 1000000001000));
assert_eq!(file.index_header.get_index_count(), 3);
assert_eq!(file.index_header.get_begin_timestamp(), 1000000000000);
assert_eq!(file.index_header.get_end_timestamp(), 1000000001000);
}
#[test]
fn test_put_key_hash_collision() {
let file = create_test_index_file("20000000000002");
let key1 = "collision_test_1";
let key2 = "collision_test_2";
file.put_key(key1, 1000, 1000000000000);
file.put_key(key2, 2000, 1000000001000);
assert_eq!(file.index_header.get_index_count(), 3); }
#[test]
fn test_is_write_full() {
use tempfile::TempDir;
let temp_dir = TempDir::new().unwrap();
let file_path = temp_dir.path().join("20000000000003");
let file_path_str = file_path.to_str().unwrap();
let file = IndexFile::new(file_path_str, 100, 5, 0, 0);
assert!(!file.is_write_full());
for i in 0..4 {
file.put_key(&format!("key{}", i), i as i64 * 1000, 1000000000000 + i as i64 * 1000);
}
assert!(file.is_write_full());
assert!(!file.put_key("overflow_key", 9999, 1000000009999));
}
#[test]
fn try_new_rejects_invalid_index_dimensions() {
use std::io::ErrorKind;
use tempfile::TempDir;
let temp_dir = TempDir::new().unwrap();
let file_path = temp_dir.path().join("20000000000030");
let file_path_str = file_path.to_str().unwrap();
let zero_slots_error = match IndexFile::try_new(file_path_str, 0, 1, 0, 0) {
Ok(_) => panic!("zero hash slots should be rejected"),
Err(error) => error,
};
assert_eq!(zero_slots_error.kind(), ErrorKind::InvalidInput);
let zero_indexes_error = match IndexFile::try_new(file_path_str, 1, 0, 0, 0) {
Ok(_) => panic!("zero index entries should be rejected"),
Err(error) => error,
};
assert_eq!(zero_indexes_error.kind(), ErrorKind::InvalidInput);
}
#[test]
fn test_time_diff_overflow_handling() {
let file = create_test_index_file("20000000000004");
file.put_key("key1", 1000, 1000000000000);
assert_eq!(file.index_header.get_begin_timestamp(), 1000000000000);
let huge_timestamp = 1000000000000 + (i32::MAX as i64 + 1000) * 1000;
file.put_key("key2", 2000, huge_timestamp);
assert_eq!(file.index_header.get_index_count(), 3);
}
#[test]
fn test_time_diff_negative_handling() {
let file = create_test_index_file("20000000000005");
file.put_key("key1", 1000, 1000000000000);
file.put_key("key2", 2000, 999999999000);
assert_eq!(file.index_header.get_index_count(), 3);
}
#[test]
fn test_select_phy_offset_basic() {
let file = create_test_index_file("20000000000006");
let begin_time = 1000000000000;
file.put_key("search_key", 12345, begin_time);
file.put_key("search_key", 23456, begin_time + 1000);
file.put_key("other_key", 99999, begin_time + 2000);
let mut results = Vec::new();
file.select_phy_offset(&mut results, "search_key", 10, begin_time - 1000, begin_time + 3000);
assert_eq!(results.len(), 2);
assert!(results.contains(&12345));
assert!(results.contains(&23456));
}
#[test]
fn test_select_phy_offset_time_range_filter() {
let file = create_test_index_file("20000000000007");
let base_time = 1000000000000;
file.put_key("key", 1000, base_time);
file.put_key("key", 2000, base_time + 5000);
file.put_key("key", 3000, base_time + 10000);
let mut results = Vec::new();
file.select_phy_offset(&mut results, "key", 10, base_time + 3000, base_time + 7000);
assert_eq!(results.len(), 1);
assert_eq!(results[0], 2000);
}
#[test]
fn test_select_phy_offset_max_num_limit() {
let file = create_test_index_file("20000000000008");
let base_time = 1000000000000;
for i in 0..10 {
file.put_key("same_key", (i * 1000) as i64, base_time + i as i64 * 1000);
}
let mut results = Vec::new();
file.select_phy_offset(&mut results, "same_key", 5, base_time - 1000, base_time + 20000);
assert_eq!(results.len(), 5, "Should respect max_num limit");
}
#[test]
fn test_is_time_matched() {
let file = create_test_index_file("20000000000009");
file.put_key("key1", 1000, 1000000000000);
file.put_key("key2", 2000, 1000000010000);
let begin_ts = file.index_header.get_begin_timestamp();
let end_ts = file.index_header.get_end_timestamp();
assert!(file.is_time_matched(begin_ts - 1000, end_ts + 1000));
assert!(file.is_time_matched(begin_ts + 1000, end_ts + 1000));
assert!(file.is_time_matched(begin_ts - 1000, end_ts - 1000));
assert!(file.is_time_matched(begin_ts + 1000, end_ts - 1000));
assert!(!file.is_time_matched(begin_ts - 5000, begin_ts - 2000));
assert!(!file.is_time_matched(end_ts + 2000, end_ts + 5000));
}
fn create_test_index_file(filename: &str) -> TestIndexFile {
use tempfile::TempDir;
let temp_dir = TempDir::new().unwrap();
let file_path = temp_dir.path().join(filename);
let file_path_str = file_path.to_str().unwrap().to_string();
let index_file = IndexFile::new(&file_path_str, 100, 1000, 0, 0);
TestIndexFile {
index_file,
_temp_dir: temp_dir, }
}
struct TestIndexFile {
index_file: IndexFile,
_temp_dir: tempfile::TempDir, }
impl std::ops::Deref for TestIndexFile {
type Target = IndexFile;
fn deref(&self) -> &Self::Target {
&self.index_file
}
}
}