///| Serializes a hint record to binary bytes.
///|
/// Layout: timestamp (8B) | expires_at (8B) | key_sz (4B) | val_sz (4B) | value_pos (8B) | key
pub fn serialize_hint_record(
timestamp : Int64,
expires_at : Int64,
key_sz : Int,
val_sz : Int,
value_pos : Int64,
key : Bytes,
) -> Bytes {
let total_sz = 32 + key_sz
let arr = Array::make(total_sz, (0).to_byte())
set_int64(arr, 0, timestamp)
set_int64(arr, 8, expires_at)
set_int32(arr, 16, key_sz)
set_int32(arr, 20, val_sz)
set_int64(arr, 24, value_pos)
for i = 0; i < key_sz; i = i + 1 {
arr[32 + i] = key[i]
}
Bytes::from_array(arr)
}
///|
/// Container for metadata read from a Hint File.
pub struct HintRecord {
timestamp : Int64
expires_at : Int64
key_sz : Int
val_sz : Int
value_pos : Int64
key : Bytes
}
///|
/// Reads and parses all hint records from a Hint File.
pub fn read_hint_records_from_file(
path : String,
) -> Array[HintRecord] raise MoonKVError {
try {
if @fs.path_exists(path) == false {
return []
}
let file_bytes = @fs.read_file_to_bytes(path)
let total_len = file_bytes.length()
let mut offset = 0
let records = Array::new()
while offset < total_len {
if total_len - offset < 32 {
raise MoonKVError::CorruptedDatabase(
"Incomplete hint record at offset \{offset} in \{path}",
)
}
let timestamp = get_int64(file_bytes, offset)
let expires_at = get_int64(file_bytes, offset + 8)
let key_sz = get_int32(file_bytes, offset + 16)
let val_sz = get_int32(file_bytes, offset + 20)
let value_pos = get_int64(file_bytes, offset + 24)
if key_sz < 0 {
raise MoonKVError::CorruptedDatabase(
"Negative key size in hint record: \{key_sz}",
)
}
let record_sz = 32 + key_sz
if offset + record_sz > total_len {
raise MoonKVError::CorruptedDatabase(
"Hint record out of bounds at offset \{offset} in \{path}",
)
}
let key_arr = Array::make(key_sz, (0).to_byte())
for i = 0; i < key_sz; i = i + 1 {
key_arr[i] = file_bytes[offset + 32 + i]
}
let key = Bytes::from_array(key_arr)
records.push({ timestamp, expires_at, key_sz, val_sz, value_pos, key })
offset = offset + record_sz
}
records
} catch {
_ =>
raise MoonKVError::IOError(
"Failed to read hint records from file \{path}",
)
}
}
///| Merges and compacts all older read-only database files.
///|
/// Writes consolidated active records to a new data and hint file, and deletes older files.
pub fn DB::merge(self : DB) -> Unit raise MoonKVError {
let files = @fs.read_dir(self.dir) catch {
_ => raise MoonKVError::IOError("Failed to read directory during merge")
}
// 1. Identify the max file ID and roll active file
let max_id = self.active_file_id
let merge_file_id = max_id + 1
self.active_file_id = max_id + 2
self.active_file_offset = 0
// Persist the next active segment immediately. Without this empty file,
// a fresh process after compaction would mistake the hint-backed merge
// segment for the active segment and append new records to it.
let active_data_path = self.dir +
"/" +
self.active_file_id.to_string() +
".data"
append_bytes_to_file(active_data_path, Bytes::make(0, (0).to_byte()))
let old_ids = Array::new()
for file in files {
if ends_with_data(file) {
let id_str = file[:file.length() - 5].to_owned()
let id = parse_int(id_str)
if id <= max_id {
old_ids.push(id)
}
}
}
if old_ids.length() == 0 {
return // No older files to merge
}
old_ids.sort()
let merge_data_path = self.dir + "/" + merge_file_id.to_string() + ".data"
let merge_hint_path = self.dir + "/" + merge_file_id.to_string() + ".hint"
let mut merge_offset = 0
// 2. Iterate through index and rewrite active records belonging to older files
let keys = self.keydir.keys()
for key in keys {
let entry = match self.keydir.get(key) {
Some(e) => e
None => continue
}
if entry.file_id <= max_id {
try {
let val = self.get(key)
let record = if entry.expires_at > 0L {
let ttl = entry.expires_at - entry.timestamp
Record::new_ttl(string_to_bytes(key), val, entry.timestamp, ttl)
} else {
Record::new(string_to_bytes(key), val, entry.timestamp)
}
let record_bytes = record.serialize()
let record_len = record_bytes.length()
// Append to consolidated merge data file
append_bytes_to_file(merge_data_path, record_bytes)
// Append to merge hint file
let hint_bytes = serialize_hint_record(
entry.timestamp,
entry.expires_at,
key.length(),
val.length(),
merge_offset.to_int64(),
string_to_bytes(key),
)
append_bytes_to_file(merge_hint_path, hint_bytes)
// Update in-memory KeyEntry to point to the new merge file location
let new_entry = KeyEntry::new(
merge_file_id,
val.length(),
merge_offset.to_int64(),
entry.timestamp,
entry.expires_at,
)
self.keydir.put(key, new_entry)
merge_offset = merge_offset + record_len
} catch {
_ => () // Skip if expired or deleted
}
}
}
// 3. Delete the older merged files
for id in old_ids {
let data_p = self.dir + "/" + id.to_string() + ".data"
let hint_p = self.dir + "/" + id.to_string() + ".hint"
try {
if @fs.path_exists(data_p) {
@fs.remove_file(data_p)
}
if @fs.path_exists(hint_p) {
@fs.remove_file(hint_p)
}
} catch {
_ => ()
}
}
}