use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use crate::idset::IdSet;
use super::Result;
pub(super) struct SparseArrayIndex {
referenced: IdSet,
mmap: memmap2::MmapMut,
_file: std::fs::File,
}
impl SparseArrayIndex {
#[allow(clippy::cast_possible_truncation)]
pub(super) fn get(&self, node_id: i64) -> Option<(i32, i32)> {
let rank = self.referenced.rank_if_set(node_id)?;
let byte_offset = (rank * 8) as usize;
if byte_offset + 8 > self.mmap.len() {
return None;
}
let packed = unsafe {
let ptr = self.mmap.as_ptr().add(byte_offset).cast::<AtomicU64>();
(*ptr).load(Ordering::Relaxed)
};
if packed == 0 {
return None;
}
let lat = packed as i32;
let lon = (packed >> 32) as i32;
Some((lat, lon))
}
}
struct SharedSparseWriter {
base: *mut u8,
capacity_bytes: usize,
}
unsafe impl Send for SharedSparseWriter {}
unsafe impl Sync for SharedSparseWriter {}
impl SharedSparseWriter {
#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
fn insert(&self, rank: u64, lat: i32, lon: i32) {
let byte_offset = (rank * 8) as usize;
if byte_offset + 8 > self.capacity_bytes {
return;
}
let packed = u64::from(lat as u32) | (u64::from(lon as u32) << 32);
unsafe {
let ptr = self.base.add(byte_offset).cast::<AtomicU64>();
(*ptr).store(packed, Ordering::Relaxed);
}
}
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
pub(super) fn build_node_index_sparse(
input: &Path,
_direct_io: bool,
scratch_dir: &Path,
mut referenced: IdSet,
) -> Result<SparseArrayIndex> {
referenced.build_rank_index();
let total_bytes = referenced.total_count().saturating_mul(8);
let temp_path = scratch_dir.join(format!(".pbfhogg-sparse-index-{}", std::process::id()));
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(&temp_path)
.map_err(|e| format!("failed to create sparse index temp file: {e}"))?;
drop(std::fs::remove_file(&temp_path));
file.set_len(total_bytes)
.map_err(|e| format!("failed to size sparse index file ({total_bytes} bytes): {e}"))?;
let mut mmap = unsafe {
memmap2::MmapMut::map_mut(&file)
.map_err(|e| format!("failed to mmap sparse index values: {e}"))?
};
let capacity_bytes = usize::try_from(total_bytes)
.map_err(|_| "sparse index total_bytes does not fit in usize")?;
let writer = SharedSparseWriter {
base: mmap.as_mut_ptr(),
capacity_bytes,
};
let (schedule, shared_input) = crate::scan::classify::build_classify_schedule(
input,
Some(crate::blob_meta::ElemKind::Node),
)?;
let referenced_ref = &referenced;
let writer_ref = &writer;
type Scratch = (Vec<crate::scan::node::NodeTuple>, Vec<(usize, usize)>);
crate::scan::classify::parallel_scan_blobs_raw(
&shared_input,
&schedule,
None,
|| -> Scratch { (Vec::new(), Vec::new()) },
|decompressed, (tuples, group_starts)| -> crate::error::Result<()> {
tuples.clear();
crate::scan::node::extract_node_tuples(decompressed, tuples, group_starts)?;
for tup in tuples.iter() {
if tup.id < 0 {
continue;
}
if let Some(rank) = referenced_ref.rank_if_set(tup.id) {
writer_ref.insert(rank, tup.lat, tup.lon);
}
}
Ok(())
},
|_seq, ()| {},
)?;
Ok(SparseArrayIndex {
referenced,
mmap,
_file: file,
})
}