oxedyne/fe2o3/fe2o3_social/src/mmap_index.rs
8.5 KiB, 1 run
created by r1870400018:9027, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! Memory-mapped index for efficient graph lookups. |
| 2 | //! |
| 3 | //! Provides O(1) access to node edges using memory-mapped files. |
| 4 | //! OS manages memory paging, keeping only accessed portions in RAM. |
| 5 | |
| 6 | use oxedyne_fe2o3_core::prelude::*; |
| 7 | |
| 8 | use std::fs::{ |
| 9 | File, |
| 10 | OpenOptions, |
| 11 | }; |
| 12 | use std::io::Write; |
| 13 | use std::path::Path; |
| 14 | |
| 15 | use memmap2::{ |
| 16 | Mmap, |
| 17 | MmapOptions, |
| 18 | }; |
| 19 | |
| 20 | /// Magic number identifying valid index files. |
| 21 | const INDEX_MAGIC: u32 = 0x4944584D; // "IDXM" in hex. |
| 22 | |
| 23 | /// Current index format version. |
| 24 | const INDEX_VERSION: u32 = 1; |
| 25 | |
| 26 | /// Header size in bytes (magic + version + node_count + reserved). |
| 27 | const HEADER_SIZE: usize = 16; |
| 28 | |
| 29 | /// Size of each directory entry in bytes (node_offset + edge_count + reserved). |
| 30 | const DIRECTORY_ENTRY_SIZE: usize = 12; |
| 31 | |
| 32 | /// Memory-mapped index for graph edge lookups. |
| 33 | /// |
| 34 | /// Format: |
| 35 | /// - Header: magic(4) + version(4) + node_count(4) + reserved(4) |
| 36 | /// - Directory: array of (data_offset(4) + edge_count(4) + reserved(4)) per node |
| 37 | /// - Data: variable-length arrays of edge_offset(8) values |
| 38 | pub struct MmapIndex { |
| 39 | /// Memory-mapped index data. |
| 40 | mmap: Mmap, |
| 41 | /// Number of nodes in the index. |
| 42 | node_count: u32, |
| 43 | } |
| 44 | |
| 45 | impl MmapIndex { |
| 46 | /// Loads an existing index file. |
| 47 | /// |
| 48 | /// # Arguments |
| 49 | /// * `path` - Path to the index file. |
| 50 | /// |
| 51 | /// # Returns |
| 52 | /// Loaded index or error if file invalid/missing. |
| 53 | pub fn load(path: &str) -> Outcome<Self> { |
| 54 | if !Path::new(path).exists() { |
| 55 | return Err(err!("Index file not found: {}", path; Missing, File)); |
| 56 | } |
| 57 | |
| 58 | let file = res!(File::open(path)); |
| 59 | let metadata = res!(file.metadata()); |
| 60 | |
| 61 | if metadata.len() < HEADER_SIZE as u64 { |
| 62 | return Err(err!("Index file too small: {}", path; Invalid, Format)); |
| 63 | } |
| 64 | |
| 65 | // Create memory mapping. |
| 66 | // SAFETY: File was opened successfully and has valid length. |
| 67 | #[allow(unsafe_code)] |
| 68 | let mmap = res!(unsafe { MmapOptions::new().map(&file) }); |
| 69 | |
| 70 | // Validate header from memory-mapped data. |
| 71 | let magic = u32::from_le_bytes([mmap[0], mmap[1], mmap[2], mmap[3]]); |
| 72 | if magic != INDEX_MAGIC { |
| 73 | return Err(err!("Invalid index magic number"; Invalid, Format)); |
| 74 | } |
| 75 | |
| 76 | let version = u32::from_le_bytes([mmap[4], mmap[5], mmap[6], mmap[7]]); |
| 77 | if version != INDEX_VERSION { |
| 78 | return Err(err!("Unsupported index version: {}", version; Invalid, Version)); |
| 79 | } |
| 80 | |
| 81 | let node_count = u32::from_le_bytes([mmap[8], mmap[9], mmap[10], mmap[11]]); |
| 82 | |
| 83 | // Validate minimum file size (header + directory). |
| 84 | let min_size = HEADER_SIZE + (node_count as usize * DIRECTORY_ENTRY_SIZE); |
| 85 | if metadata.len() < min_size as u64 { |
| 86 | return Err(err!("Index file too small: expected at least {} bytes", min_size; Invalid, Format)); |
| 87 | } |
| 88 | |
| 89 | Ok(Self { mmap, node_count }) |
| 90 | } |
| 91 | |
| 92 | /// Gets all edge offsets for a node. |
| 93 | /// |
| 94 | /// # Arguments |
| 95 | /// * `node_id` - The node ID to look up. |
| 96 | /// |
| 97 | /// # Returns |
| 98 | /// Vector of edge offsets in the edge file, or None if node doesn't exist. |
| 99 | pub fn get_node_edge_offsets(&self, node_id: u32) -> Option<Vec<u64>> { |
| 100 | if node_id >= self.node_count { |
| 101 | return None; |
| 102 | } |
| 103 | |
| 104 | // Calculate offset to this node's directory entry. |
| 105 | let dir_offset = HEADER_SIZE + (node_id as usize * DIRECTORY_ENTRY_SIZE); |
| 106 | |
| 107 | if dir_offset + DIRECTORY_ENTRY_SIZE > self.mmap.len() { |
| 108 | return None; |
| 109 | } |
| 110 | |
| 111 | // Read directory entry: data_offset(4) + edge_count(4) + reserved(4). |
| 112 | let data_offset = u32::from_le_bytes([ |
| 113 | self.mmap[dir_offset], self.mmap[dir_offset + 1], |
| 114 | self.mmap[dir_offset + 2], self.mmap[dir_offset + 3], |
| 115 | ]) as usize; |
| 116 | |
| 117 | let edge_count = u32::from_le_bytes([ |
| 118 | self.mmap[dir_offset + 4], self.mmap[dir_offset + 5], |
| 119 | self.mmap[dir_offset + 6], self.mmap[dir_offset + 7], |
| 120 | ]) as usize; |
| 121 | |
| 122 | // If no edges, return empty vector. |
| 123 | if edge_count == 0 { |
| 124 | return Some(Vec::new()); |
| 125 | } |
| 126 | |
| 127 | // Validate data bounds. |
| 128 | let data_end = data_offset + (edge_count * 8); // 8 bytes per u64 offset |
| 129 | if data_end > self.mmap.len() { |
| 130 | return None; |
| 131 | } |
| 132 | |
| 133 | // Read edge offsets from data section. |
| 134 | let mut edge_offsets = Vec::with_capacity(edge_count); |
| 135 | for i in 0..edge_count { |
| 136 | let offset_pos = data_offset + (i * 8); |
| 137 | let edge_offset = u64::from_le_bytes([ |
| 138 | self.mmap[offset_pos], self.mmap[offset_pos + 1], |
| 139 | self.mmap[offset_pos + 2], self.mmap[offset_pos + 3], |
| 140 | self.mmap[offset_pos + 4], self.mmap[offset_pos + 5], |
| 141 | self.mmap[offset_pos + 6], self.mmap[offset_pos + 7], |
| 142 | ]); |
| 143 | edge_offsets.push(edge_offset); |
| 144 | } |
| 145 | |
| 146 | Some(edge_offsets) |
| 147 | } |
| 148 | |
| 149 | /// Gets the total number of nodes. |
| 150 | pub fn node_count(&self) -> u32 { |
| 151 | self.node_count |
| 152 | } |
| 153 | } |
| 154 | |
| 155 | /// Builder for creating memory-mapped indices during graph generation. |
| 156 | pub struct MmapIndexBuilder { |
| 157 | /// Path for the index file. |
| 158 | file_path: String, |
| 159 | /// Number of nodes to index. |
| 160 | node_count: u32, |
| 161 | /// Tracks edge offsets per node. |
| 162 | node_edges: Vec<Vec<u64>>, |
| 163 | } |
| 164 | |
| 165 | impl MmapIndexBuilder { |
| 166 | /// Creates a new index builder. |
| 167 | /// |
| 168 | /// # Arguments |
| 169 | /// * `path` - Path for the index file. |
| 170 | /// * `node_count` - Total number of nodes to index. |
| 171 | /// |
| 172 | /// # Returns |
| 173 | /// New builder instance. |
| 174 | pub fn new(path: &str, node_count: u32) -> Outcome<Self> { |
| 175 | Ok(Self { |
| 176 | file_path: path.to_string(), |
| 177 | node_count, |
| 178 | node_edges: vec![Vec::new(); node_count as usize], |
| 179 | }) |
| 180 | } |
| 181 | |
| 182 | /// Records an edge for a node. |
| 183 | /// |
| 184 | /// # Arguments |
| 185 | /// * `from_node` - Source node ID. |
| 186 | /// * `edge_offset` - Offset of this edge in the edge file. |
| 187 | pub fn add_edge(&mut self, from_node: u32, edge_offset: u64) -> Outcome<()> { |
| 188 | if from_node >= self.node_count { |
| 189 | return Err(err!("Node ID {} exceeds node count {}", from_node, self.node_count; Invalid, Input)); |
| 190 | } |
| 191 | |
| 192 | self.node_edges[from_node as usize].push(edge_offset); |
| 193 | Ok(()) |
| 194 | } |
| 195 | |
| 196 | /// Finalises the index by writing all data to disk. |
| 197 | /// |
| 198 | /// Must be called after all edges have been added. |
| 199 | pub fn finalise(self) -> Outcome<()> { |
| 200 | let mut file = res!(OpenOptions::new() |
| 201 | .read(true) |
| 202 | .write(true) |
| 203 | .create(true) |
| 204 | .truncate(true) |
| 205 | .open(&self.file_path)); |
| 206 | |
| 207 | // Calculate file size. |
| 208 | let directory_size = self.node_count as usize * DIRECTORY_ENTRY_SIZE; |
| 209 | let total_edges: usize = self.node_edges.iter().map(|edges| edges.len()).sum(); |
| 210 | let data_size = total_edges * 8; // 8 bytes per u64 offset |
| 211 | let file_size = HEADER_SIZE + directory_size + data_size; |
| 212 | |
| 213 | // Pre-allocate file space. |
| 214 | res!(file.set_len(file_size as u64)); |
| 215 | |
| 216 | // Write header. |
| 217 | res!(file.write_all(&INDEX_MAGIC.to_le_bytes())); |
| 218 | res!(file.write_all(&INDEX_VERSION.to_le_bytes())); |
| 219 | res!(file.write_all(&self.node_count.to_le_bytes())); |
| 220 | res!(file.write_all(&[0u8; 4])); // Reserved. |
| 221 | |
| 222 | // Calculate data section start. |
| 223 | let data_start = HEADER_SIZE + directory_size; |
| 224 | let mut current_data_offset = data_start; |
| 225 | |
| 226 | // Write directory entries and collect data to write. |
| 227 | let mut data_buffer = Vec::with_capacity(data_size); |
| 228 | |
| 229 | for node_id in 0..self.node_count { |
| 230 | let edges = &self.node_edges[node_id as usize]; |
| 231 | let edge_count = edges.len() as u32; |
| 232 | |
| 233 | // Write directory entry: data_offset(4) + edge_count(4) + reserved(4). |
| 234 | res!(file.write_all(&(current_data_offset as u32).to_le_bytes())); |
| 235 | res!(file.write_all(&edge_count.to_le_bytes())); |
| 236 | res!(file.write_all(&[0u8; 4])); // Reserved. |
| 237 | |
| 238 | // Add edge offsets to data buffer. |
| 239 | for &edge_offset in edges { |
| 240 | data_buffer.extend_from_slice(&edge_offset.to_le_bytes()); |
| 241 | } |
| 242 | |
| 243 | // Update data offset for next node. |
| 244 | current_data_offset += edges.len() * 8; |
| 245 | } |
| 246 | |
| 247 | // Write all edge offset data. |
| 248 | res!(file.write_all(&data_buffer)); |
| 249 | res!(file.flush()); |
| 250 | |
| 251 | Ok(()) |
| 252 | } |
| 253 | } |