oxedyne/fe2o3/fe2o3_social/src/mmap_graph.rs
24.3 KiB, 126 runs
created by r1870400018:8969, 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 graph storage for large social networks. |
| 2 | //! |
| 3 | //! Provides disk-backed storage for graphs that exceed available RAM, |
| 4 | //! using memory-mapped files for efficient access patterns. |
| 5 | |
| 6 | use crate::{ |
| 7 | graph::GraphAccessMethod, |
| 8 | mmap_index::{ |
| 9 | MmapIndex, |
| 10 | MmapIndexBuilder, |
| 11 | }, |
| 12 | }; |
| 13 | |
| 14 | use oxedyne_fe2o3_core::{ |
| 15 | prelude::*, |
| 16 | info, |
| 17 | warn, |
| 18 | }; |
| 19 | |
| 20 | use std::{ |
| 21 | cell::RefCell, |
| 22 | fs::{ |
| 23 | File, |
| 24 | OpenOptions, |
| 25 | }, |
| 26 | io::{ |
| 27 | Write, |
| 28 | Read, |
| 29 | Seek, |
| 30 | SeekFrom, |
| 31 | }, |
| 32 | }; |
| 33 | |
| 34 | #[cfg(unix)] |
| 35 | use std::os::unix::io::AsRawFd; |
| 36 | |
| 37 | use memmap2::{ |
| 38 | Mmap, |
| 39 | MmapOptions, |
| 40 | }; |
| 41 | |
| 42 | |
| 43 | /// Memory-mapped graph storage using flat edge list. |
| 44 | /// |
| 45 | /// Stores edges in binary format on disk, memory-maps for access. |
| 46 | /// Each edge is 9 bytes: from_id(4) + to_id(4) + link_data(1). |
| 47 | pub struct MmapGraph { |
| 48 | /// File handle (kept for writing during generation). |
| 49 | edge_file: File, |
| 50 | /// Number of edges stored. |
| 51 | num_edges: usize, |
| 52 | /// Maximum edges capacity. |
| 53 | capacity: usize, |
| 54 | /// Memory-mapped view of the edge data (for fast reading). |
| 55 | mmap: Option<Mmap>, |
| 56 | /// Disk-based index for fast lookups. |
| 57 | /// Loaded from .idx file if available. |
| 58 | index: Option<MmapIndex>, |
| 59 | /// File path for the edge data. |
| 60 | file_path: String, |
| 61 | /// Opening and closing the file can degrade performance. |
| 62 | read_handle: RefCell<Option<File>>, |
| 63 | } |
| 64 | |
| 65 | impl MmapGraph { |
| 66 | /// Creates a new memory-mapped graph with indexing. |
| 67 | /// |
| 68 | /// # Arguments |
| 69 | /// * `path` - Path for the edge data file. |
| 70 | /// * `capacity` - Maximum number of edges to support. |
| 71 | /// |
| 72 | /// # Returns |
| 73 | /// New memory-mapped graph instance. |
| 74 | pub fn new(path: &str, capacity: usize) -> Outcome<Self> { |
| 75 | Self::new_with_options(path, capacity, true) |
| 76 | } |
| 77 | |
| 78 | /// Creates a new memory-mapped graph without indexing for large graphs. |
| 79 | /// |
| 80 | /// Disables the in-memory index to save memory for very large graphs. |
| 81 | /// Graph queries will be slower but memory usage will be minimal. |
| 82 | /// |
| 83 | /// # Arguments |
| 84 | /// * `path` - Path for the edge data file. |
| 85 | /// * `capacity` - Maximum number of edges to support. |
| 86 | /// |
| 87 | /// # Returns |
| 88 | /// New memory-mapped graph instance without indexing. |
| 89 | pub fn new_without_index(path: &str, capacity: usize) -> Outcome<Self> { |
| 90 | Self::new_with_options(path, capacity, false) |
| 91 | } |
| 92 | |
| 93 | /// Creates a new memory-mapped graph with optional indexing. |
| 94 | /// |
| 95 | /// # Arguments |
| 96 | /// * `path` - Path for the edge data file. |
| 97 | /// * `capacity` - Maximum number of edges to support. |
| 98 | /// * `enable_index` - Whether to enable in-memory indexing. |
| 99 | /// |
| 100 | /// # Returns |
| 101 | /// New memory-mapped graph instance. |
| 102 | pub fn new_with_options(path: &str, capacity: usize, _enable_index: bool) -> Outcome<Self> { |
| 103 | // Create or truncate the file. |
| 104 | let edge_file = res!(OpenOptions::new() |
| 105 | .read(true) |
| 106 | .write(true) |
| 107 | .create(true) |
| 108 | .truncate(true) |
| 109 | .open(path)); |
| 110 | |
| 111 | // Pre-allocate space for edges. |
| 112 | let file_size = capacity * 9; // 9 bytes per edge. |
| 113 | res!(edge_file.set_len(file_size as u64)); |
| 114 | |
| 115 | Ok(Self { |
| 116 | edge_file, |
| 117 | num_edges: 0, |
| 118 | capacity, |
| 119 | mmap: None, // Memory mapping created when loading existing graphs. |
| 120 | index: None, // Index will be loaded separately after generation. |
| 121 | file_path: path.to_string(), |
| 122 | read_handle: RefCell::new(None), |
| 123 | }) |
| 124 | } |
| 125 | |
| 126 | /// Loads an existing memory-mapped graph. |
| 127 | /// |
| 128 | /// # Arguments |
| 129 | /// * `path` - Path to existing edge data file. |
| 130 | /// |
| 131 | /// # Returns |
| 132 | /// Loaded memory-mapped graph instance. |
| 133 | pub fn load_existing(path: &str) -> Outcome<Self> { |
| 134 | // Open existing file without truncating. |
| 135 | let edge_file = res!(OpenOptions::new() |
| 136 | .read(true) |
| 137 | .write(true) |
| 138 | .open(path)); |
| 139 | |
| 140 | // Get file size to determine capacity and edge count. |
| 141 | let metadata = res!(edge_file.metadata()); |
| 142 | let file_size = metadata.len(); |
| 143 | let capacity = (file_size / 9) as usize; // 9 bytes per edge. |
| 144 | let num_edges = capacity; // Assume file is fully written for now. |
| 145 | |
| 146 | // Create memory mapping for fast reading. |
| 147 | // SAFETY: The file was just opened successfully and has valid content. |
| 148 | // We only read from this mapping, never write to it. |
| 149 | #[allow(unsafe_code)] |
| 150 | let mmap = unsafe { |
| 151 | MmapOptions::new().map(&edge_file) |
| 152 | }; |
| 153 | |
| 154 | let mmap = match mmap { |
| 155 | Ok(m) => { |
| 156 | info!("Created memory mapping for edge data ({:.1} MB)", m.len() as f64 / 1024.0 / 1024.0); |
| 157 | // Tell OS we don't need aggressive caching. |
| 158 | #[cfg(unix)] |
| 159 | #[allow(unsafe_code)] |
| 160 | unsafe { |
| 161 | let result = libc::madvise( |
| 162 | m.as_ptr() as *mut libc::c_void, |
| 163 | m.len(), |
| 164 | libc::MADV_RANDOM // Random access pattern, don't prefetch |
| 165 | ); |
| 166 | if result == 0 { |
| 167 | info!("Applied MADV_RANDOM to memory mapping"); |
| 168 | } else { |
| 169 | warn!("Failed to apply madvise: {}", result); |
| 170 | } |
| 171 | } |
| 172 | Some(m) |
| 173 | }, |
| 174 | Err(e) => { |
| 175 | warn!("Failed to create memory mapping: {}. Will use direct file access.", e); |
| 176 | None |
| 177 | } |
| 178 | }; |
| 179 | |
| 180 | let mut graph = Self { |
| 181 | edge_file, |
| 182 | num_edges, |
| 183 | capacity, |
| 184 | mmap, |
| 185 | index: None, |
| 186 | file_path: path.to_string(), |
| 187 | read_handle: RefCell::new(None), |
| 188 | }; |
| 189 | |
| 190 | // Try to load the disk-based index if it exists. |
| 191 | let index_path = format!("{}.idx", path.trim_end_matches(".mmap")); |
| 192 | match MmapIndex::load(&index_path) { |
| 193 | Ok(idx) => { |
| 194 | info!("Loaded disk-based index for fast lookups"); |
| 195 | graph.index = Some(idx); |
| 196 | }, |
| 197 | Err(_) => { |
| 198 | info!("No index found for {}. Graph lookups will scan edges.", path); |
| 199 | } |
| 200 | } |
| 201 | |
| 202 | Ok(graph) |
| 203 | } |
| 204 | |
| 205 | pub fn release_memory(&self) { |
| 206 | #[cfg(unix)] |
| 207 | if let Some(ref mmap) = self.mmap { |
| 208 | #[allow(unsafe_code)] |
| 209 | unsafe { |
| 210 | // This tells the OS it can free these pages. |
| 211 | // They'll be reloaded on next access. |
| 212 | let result = libc::madvise( |
| 213 | mmap.as_ptr() as *mut libc::c_void, |
| 214 | mmap.len(), |
| 215 | libc::MADV_DONTNEED // Free the pages. |
| 216 | ); |
| 217 | if result == 0 { |
| 218 | debug!("Released mmap memory pages"); |
| 219 | } |
| 220 | } |
| 221 | } |
| 222 | } |
| 223 | |
| 224 | /// Adds an edge to the graph. |
| 225 | /// |
| 226 | /// # Arguments |
| 227 | /// * `from` - Source node ID. |
| 228 | /// * `to` - Target node ID. |
| 229 | /// * `link_data` - Packed link data (1 byte). |
| 230 | /// |
| 231 | /// # Returns |
| 232 | /// Ok if successful, error if capacity exceeded. |
| 233 | pub fn add_edge(&mut self, from: u32, to: u32, link_data: u8) -> Outcome<()> { |
| 234 | if self.num_edges >= self.capacity { |
| 235 | return Err(err!("Graph capacity exceeded: {} edges", self.capacity; Invalid, Input)); |
| 236 | } |
| 237 | |
| 238 | // Index tracking handled by the builder during generation. |
| 239 | |
| 240 | // Write edge data directly to file. |
| 241 | res!(self.edge_file.write_all(&from.to_le_bytes())); |
| 242 | res!(self.edge_file.write_all(&to.to_le_bytes())); |
| 243 | res!(self.edge_file.write_all(&[link_data])); |
| 244 | |
| 245 | self.num_edges += 1; |
| 246 | Ok(()) |
| 247 | } |
| 248 | |
| 249 | /// Flushes pending writes to disk. |
| 250 | pub fn flush(&mut self) -> Outcome<()> { |
| 251 | res!(self.edge_file.flush()); |
| 252 | Ok(()) |
| 253 | } |
| 254 | |
| 255 | |
| 256 | |
| 257 | /// Gets the number of edges stored. |
| 258 | pub fn edge_count(&self) -> usize { |
| 259 | self.num_edges |
| 260 | } |
| 261 | |
| 262 | /// Gets memory usage in MB (always near zero for mmap). |
| 263 | pub fn memory_usage_mb(&self) -> f64 { |
| 264 | // Only the file handle and metadata are in memory. |
| 265 | // Actual edge data is on disk. |
| 266 | 0.001 // Negligible memory usage. |
| 267 | } |
| 268 | |
| 269 | /// Gets incoming edges to a node ID. |
| 270 | /// Scans all edges to find matches (O(n) operation). |
| 271 | /// Note: This is expensive for large graphs without reverse index. |
| 272 | pub fn get_incoming_edges(&self, to_id: u32) -> Outcome<Vec<(u32, u8)>> { |
| 273 | let mut edges = Vec::new(); |
| 274 | |
| 275 | if let Some(ref mmap) = self.mmap { |
| 276 | // Scan all edges to find incoming ones. |
| 277 | let edge_size = 9; // 4 bytes from_id + 4 bytes to_id + 1 byte link_data. |
| 278 | let num_edges = self.num_edges; |
| 279 | |
| 280 | for i in 0..num_edges { |
| 281 | let offset = i * edge_size; |
| 282 | if offset + edge_size > mmap.len() { |
| 283 | break; |
| 284 | } |
| 285 | |
| 286 | // Read edge. |
| 287 | let edge_from = u32::from_le_bytes([ |
| 288 | mmap[offset], mmap[offset + 1], mmap[offset + 2], mmap[offset + 3] |
| 289 | ]); |
| 290 | let edge_to = u32::from_le_bytes([ |
| 291 | mmap[offset + 4], mmap[offset + 5], mmap[offset + 6], mmap[offset + 7] |
| 292 | ]); |
| 293 | let link_data = mmap[offset + 8]; |
| 294 | |
| 295 | // Check if this edge points to our target node. |
| 296 | if edge_to == to_id { |
| 297 | edges.push((edge_from, link_data)); |
| 298 | } |
| 299 | } |
| 300 | |
| 301 | return Ok(edges); |
| 302 | } |
| 303 | |
| 304 | // Fallback: read from file if memory mapping is unavailable. |
| 305 | warn!("Memory mapping unavailable for incoming edges to node {}, using direct file access", to_id); |
| 306 | |
| 307 | let mut file = res!(File::open(&self.file_path)); |
| 308 | let edge_size = 9; |
| 309 | |
| 310 | for i in 0..self.num_edges { |
| 311 | let offset = i as u64 * edge_size as u64; |
| 312 | res!(file.seek(SeekFrom::Start(offset))); |
| 313 | |
| 314 | let mut buffer = [0u8; 9]; |
| 315 | match file.read_exact(&mut buffer) { |
| 316 | Ok(()) => { |
| 317 | let edge_from = u32::from_le_bytes([buffer[0], buffer[1], buffer[2], buffer[3]]); |
| 318 | let edge_to = u32::from_le_bytes([buffer[4], buffer[5], buffer[6], buffer[7]]); |
| 319 | let link_data = buffer[8]; |
| 320 | |
| 321 | if edge_to == to_id { |
| 322 | edges.push((edge_from, link_data)); |
| 323 | } |
| 324 | } |
| 325 | Err(_) => break, // End of file or read error. |
| 326 | } |
| 327 | } |
| 328 | |
| 329 | Ok(edges) |
| 330 | } |
| 331 | |
| 332 | /// Retrieves all outgoing edges from a specified node. |
| 333 | /// |
| 334 | /// This method supports multiple access strategies to balance memory usage and performance: |
| 335 | /// - FileIO: Always uses direct file I/O (lowest memory, slowest). |
| 336 | /// - Mmap: Always uses memory mapping (highest memory, fastest). |
| 337 | /// - Auto: Chooses based on edge count threshold. |
| 338 | /// |
| 339 | /// # Arguments |
| 340 | /// * `from_id` - The source node ID to retrieve edges from. |
| 341 | /// * `method` - The access method to use for retrieving edges. |
| 342 | /// |
| 343 | /// # Returns |
| 344 | /// A vector of tuples containing (target_node_id, link_data) for all outgoing edges. |
| 345 | pub fn get_outgoing_edges( |
| 346 | &self, |
| 347 | from_id: u32, |
| 348 | method: GraphAccessMethod, |
| 349 | ) |
| 350 | -> Outcome<Vec<(u32, u8)>> |
| 351 | { |
| 352 | let mut edges = Vec::new(); |
| 353 | |
| 354 | // Try index-based lookup first. |
| 355 | if let Some(ref index) = self.index { |
| 356 | if let Some(edge_offsets) = index.get_node_edge_offsets(from_id) { |
| 357 | if edge_offsets.is_empty() { |
| 358 | return Ok(edges); |
| 359 | } |
| 360 | |
| 361 | // Decide which method to use. |
| 362 | let use_file_io = match method { |
| 363 | GraphAccessMethod::FileIO => true, |
| 364 | GraphAccessMethod::Mmap => false, |
| 365 | GraphAccessMethod::Auto(edge_lim) => edge_offsets.len() < edge_lim, |
| 366 | }; |
| 367 | |
| 368 | if use_file_io { |
| 369 | // Get or create cached file handle. |
| 370 | let mut read_handle = self.read_handle.borrow_mut(); |
| 371 | if read_handle.is_none() { |
| 372 | *read_handle = Some(res!(File::open(&self.file_path))); |
| 373 | } |
| 374 | |
| 375 | let file = match read_handle.as_mut() { |
| 376 | Some(f) => f, |
| 377 | None => return Err(err!("Failed to get file handle after creation"; Bug)), |
| 378 | }; |
| 379 | |
| 380 | // Find the range of bytes we need. |
| 381 | let first_offset = match edge_offsets.first() { |
| 382 | Some(offset) => *offset, |
| 383 | None => return Err(err!("Edge offsets unexpectedly empty"; Invalid, Input)), |
| 384 | }; |
| 385 | let last_offset = match edge_offsets.last() { |
| 386 | Some(offset) => *offset, |
| 387 | None => return Err(err!("Edge offsets unexpectedly empty"; Invalid, Input)), |
| 388 | }; |
| 389 | let bytes_to_read = (last_offset - first_offset + 9) as usize; |
| 390 | |
| 391 | // Read all relevant bytes in one operation. |
| 392 | res!(file.seek(SeekFrom::Start(first_offset))); |
| 393 | let mut buffer = vec![0u8; bytes_to_read]; |
| 394 | res!(file.read_exact(&mut buffer)); |
| 395 | |
| 396 | // Parse edges from buffer. |
| 397 | for edge_offset in edge_offsets { |
| 398 | let relative_offset = (edge_offset - first_offset) as usize; |
| 399 | if relative_offset + 9 <= buffer.len() { |
| 400 | let edge_from = u32::from_le_bytes([ |
| 401 | buffer[relative_offset], |
| 402 | buffer[relative_offset + 1], |
| 403 | buffer[relative_offset + 2], |
| 404 | buffer[relative_offset + 3] |
| 405 | ]); |
| 406 | let edge_to = u32::from_le_bytes([ |
| 407 | buffer[relative_offset + 4], |
| 408 | buffer[relative_offset + 5], |
| 409 | buffer[relative_offset + 6], |
| 410 | buffer[relative_offset + 7] |
| 411 | ]); |
| 412 | let link_data = buffer[relative_offset + 8]; |
| 413 | |
| 414 | if edge_from == from_id { |
| 415 | edges.push((edge_to, link_data)); |
| 416 | } |
| 417 | } |
| 418 | } |
| 419 | |
| 420 | return Ok(edges); |
| 421 | } else { |
| 422 | // Memory-mapped path for better performance. |
| 423 | if let Some(ref mmap) = self.mmap { |
| 424 | for edge_offset in edge_offsets { |
| 425 | let offset = edge_offset as usize; |
| 426 | if offset + 9 > mmap.len() { |
| 427 | continue; |
| 428 | } |
| 429 | let edge_from = u32::from_le_bytes([ |
| 430 | mmap[offset], mmap[offset + 1], mmap[offset + 2], mmap[offset + 3] |
| 431 | ]); |
| 432 | let edge_to = u32::from_le_bytes([ |
| 433 | mmap[offset + 4], mmap[offset + 5], mmap[offset + 6], mmap[offset + 7] |
| 434 | ]); |
| 435 | let link_data = mmap[offset + 8]; |
| 436 | |
| 437 | if edge_from == from_id { |
| 438 | edges.push((edge_to, link_data)); |
| 439 | } |
| 440 | } |
| 441 | return Ok(edges); |
| 442 | } |
| 443 | } |
| 444 | } |
| 445 | } |
| 446 | |
| 447 | // Fallback to mmap scanning when no index available. |
| 448 | if let Some(ref mmap) = self.mmap { |
| 449 | let mut offset = 0; |
| 450 | for _ in 0..self.num_edges { |
| 451 | if offset + 9 > mmap.len() { |
| 452 | break; |
| 453 | } |
| 454 | let edge_from = u32::from_le_bytes([ |
| 455 | mmap[offset], mmap[offset + 1], mmap[offset + 2], mmap[offset + 3] |
| 456 | ]); |
| 457 | let edge_to = u32::from_le_bytes([ |
| 458 | mmap[offset + 4], mmap[offset + 5], mmap[offset + 6], mmap[offset + 7] |
| 459 | ]); |
| 460 | let link_data = mmap[offset + 8]; |
| 461 | |
| 462 | if edge_from == from_id { |
| 463 | edges.push((edge_to, link_data)); |
| 464 | } |
| 465 | offset += 9; |
| 466 | } |
| 467 | } else { |
| 468 | // No mmap available, use direct file access. |
| 469 | warn!("Memory mapping unavailable for node {}, using direct file access", from_id); |
| 470 | #[cfg(unix)] |
| 471 | let mut file = res!(std::fs::File::open(format!("/proc/self/fd/{}", self.edge_file.as_raw_fd()))); |
| 472 | #[cfg(not(unix))] |
| 473 | let mut file = res!(std::fs::File::open(&self.file_path)); |
| 474 | |
| 475 | res!(file.seek(SeekFrom::Start(0))); |
| 476 | let mut buffer = [0u8; 9]; |
| 477 | for _ in 0..self.num_edges { |
| 478 | match file.read_exact(&mut buffer) { |
| 479 | Ok(_) => { |
| 480 | let edge_from = u32::from_le_bytes([buffer[0], buffer[1], buffer[2], buffer[3]]); |
| 481 | let edge_to = u32::from_le_bytes([buffer[4], buffer[5], buffer[6], buffer[7]]); |
| 482 | let link_data = buffer[8]; |
| 483 | if edge_from == from_id { |
| 484 | edges.push((edge_to, link_data)); |
| 485 | } |
| 486 | }, |
| 487 | Err(_) => break, |
| 488 | } |
| 489 | } |
| 490 | } |
| 491 | |
| 492 | Ok(edges) |
| 493 | } |
| 494 | |
| 495 | /// Gets the total number of edges for statistics. |
| 496 | pub fn total_edges(&self) -> usize { |
| 497 | self.num_edges |
| 498 | } |
| 499 | |
| 500 | /// Loads the memory-mapped index if available. |
| 501 | pub fn load_index(&mut self, mmap_path: &str) -> Outcome<()> { |
| 502 | let index_path = format!("{}.idx", mmap_path.trim_end_matches(".mmap")); |
| 503 | match MmapIndex::load(&index_path) { |
| 504 | Ok(idx) => { |
| 505 | self.index = Some(idx); |
| 506 | Ok(()) |
| 507 | }, |
| 508 | Err(e) => { |
| 509 | warn!("Could not load index: {}", e); |
| 510 | Ok(()) // Not having an index is ok, just slower. |
| 511 | } |
| 512 | } |
| 513 | } |
| 514 | |
| 515 | /// Builds a memory-mapped index for an existing graph by scanning all edges. |
| 516 | /// |
| 517 | /// This creates a fast lookup index for graphs that were generated without indexing. |
| 518 | /// Call this once after graph generation to enable O(1) edge lookups. |
| 519 | pub fn build_index(&mut self, mmap_path: &str, max_node_id: u32) -> Outcome<()> { |
| 520 | info!("Building memory-mapped index by scanning {} edges...", self.num_edges); |
| 521 | |
| 522 | let index_path = format!("{}.idx", mmap_path.trim_end_matches(".mmap")); |
| 523 | let mut index_builder = res!(MmapIndexBuilder::new(&index_path, max_node_id + 1)); |
| 524 | |
| 525 | // Scan all edges to build the index. |
| 526 | if let Some(ref mmap) = self.mmap { |
| 527 | // Use memory-mapped edge data for scanning. |
| 528 | let mut offset = 0; |
| 529 | for edge_idx in 0..self.num_edges { |
| 530 | if offset + 9 > mmap.len() { |
| 531 | break; |
| 532 | } |
| 533 | |
| 534 | // Read the source node ID from this edge. |
| 535 | let from_node = u32::from_le_bytes([ |
| 536 | mmap[offset], mmap[offset + 1], mmap[offset + 2], mmap[offset + 3] |
| 537 | ]); |
| 538 | |
| 539 | // Record this edge's offset in the index. |
| 540 | let edge_offset = (edge_idx * 9) as u64; |
| 541 | res!(index_builder.add_edge(from_node, edge_offset)); |
| 542 | |
| 543 | offset += 9; |
| 544 | } |
| 545 | } else { |
| 546 | return Err(err!("Cannot build index: no memory mapping available"; Invalid, Input)); |
| 547 | } |
| 548 | |
| 549 | // Finalise and save the index. |
| 550 | res!(index_builder.finalise()); |
| 551 | info!("Index built successfully: {}", index_path); |
| 552 | |
| 553 | // Load the newly created index. |
| 554 | res!(self.load_index(mmap_path)); |
| 555 | info!("Memory-mapped index loaded for fast lookups"); |
| 556 | |
| 557 | Ok(()) |
| 558 | } |
| 559 | } |
| 560 | |
| 561 | /// Builder for creating memory-mapped graphs from stub matching. |
| 562 | pub struct MmapGraphBuilder { |
| 563 | graph: MmapGraph, |
| 564 | /// Optional index builder for creating disk-based index. |
| 565 | index_builder: Option<MmapIndexBuilder>, |
| 566 | /// Path for the graph files. |
| 567 | base_path: String, |
| 568 | /// Maximum node ID for index sizing. |
| 569 | _max_node_id: u32, |
| 570 | } |
| 571 | |
| 572 | impl MmapGraphBuilder { |
| 573 | /// Creates a new builder with disk-based indexing. |
| 574 | /// |
| 575 | /// # Arguments |
| 576 | /// * `path` - Path for the edge data file. |
| 577 | /// * `estimated_edges` - Estimated number of edges (for pre-allocation). |
| 578 | /// * `max_node_id` - Maximum node ID (for index sizing). |
| 579 | /// |
| 580 | /// # Returns |
| 581 | /// New builder instance. |
| 582 | pub fn new(path: &str, estimated_edges: usize, max_node_id: u32) -> Outcome<Self> { |
| 583 | let index_path = format!("{}.idx", path.trim_end_matches(".mmap")); |
| 584 | let index_builder = Some(res!(MmapIndexBuilder::new(&index_path, max_node_id + 1))); |
| 585 | |
| 586 | Ok(Self { |
| 587 | graph: res!(MmapGraph::new_with_options(path, estimated_edges, false)), |
| 588 | index_builder, |
| 589 | base_path: path.to_string(), |
| 590 | _max_node_id: max_node_id, |
| 591 | }) |
| 592 | } |
| 593 | |
| 594 | /// Creates a new builder without indexing for large graphs. |
| 595 | /// |
| 596 | /// Disables indexing completely to save memory for very large graphs. |
| 597 | /// Graph queries will be slower but memory usage will be minimal. |
| 598 | /// |
| 599 | /// # Arguments |
| 600 | /// * `path` - Path for the edge data file. |
| 601 | /// * `estimated_edges` - Estimated number of edges (for pre-allocation). |
| 602 | /// |
| 603 | /// # Returns |
| 604 | /// New builder instance without indexing. |
| 605 | pub fn new_without_index(path: &str, estimated_edges: usize) -> Outcome<Self> { |
| 606 | Ok(Self { |
| 607 | graph: res!(MmapGraph::new_with_options(path, estimated_edges, false)), |
| 608 | index_builder: None, |
| 609 | base_path: path.to_string(), |
| 610 | _max_node_id: 0, |
| 611 | }) |
| 612 | } |
| 613 | |
| 614 | /// Adds an edge directly to the graph. |
| 615 | /// |
| 616 | /// # Arguments |
| 617 | /// * `from` - Source node ID. |
| 618 | /// * `to` - Target node ID. |
| 619 | /// * `link_data` - Edge data byte. |
| 620 | /// |
| 621 | /// # Returns |
| 622 | /// Success result. |
| 623 | pub fn add_edge(&mut self, from: u32, to: u32, link_data: u8) -> Outcome<()> { |
| 624 | // Record edge offset in index if enabled. |
| 625 | if let Some(ref mut index_builder) = self.index_builder { |
| 626 | let edge_offset = (self.graph.num_edges * 9) as u64; |
| 627 | res!(index_builder.add_edge(from, edge_offset)); |
| 628 | } |
| 629 | |
| 630 | res!(self.graph.add_edge(from, to, link_data)); |
| 631 | Ok(()) |
| 632 | } |
| 633 | |
| 634 | /// Adds an edge if it doesn't already exist. |
| 635 | /// |
| 636 | /// # Arguments |
| 637 | /// * `from` - Source node ID. |
| 638 | /// * `to` - Target node ID. |
| 639 | /// * `link_data` - Packed link data. |
| 640 | /// |
| 641 | /// # Returns |
| 642 | /// True if edge was added, false if it already existed. |
| 643 | pub fn add_edge_unique(&mut self, from: u32, to: u32, link_data: u8) -> Outcome<bool> { |
| 644 | // Note: This method is deprecated for large graphs due to memory usage. |
| 645 | // Use add_edge() directly for better memory efficiency. |
| 646 | res!(self.graph.add_edge(from, to, link_data)); |
| 647 | Ok(true) |
| 648 | } |
| 649 | |
| 650 | /// Finalises the graph and returns it. |
| 651 | /// |
| 652 | /// Writes the disk-based index if indexing was enabled. |
| 653 | pub fn build(mut self) -> Outcome<MmapGraph> { |
| 654 | res!(self.graph.flush()); |
| 655 | |
| 656 | // Write the disk-based index if we have one. |
| 657 | if let Some(index_builder) = self.index_builder { |
| 658 | res!(index_builder.finalise()); |
| 659 | info!("Disk-based index created for fast graph lookups"); |
| 660 | |
| 661 | // Load the index into the graph. |
| 662 | res!(self.graph.load_index(&self.base_path)); |
| 663 | } else { |
| 664 | info!("No index created (indexing was disabled)"); |
| 665 | } |
| 666 | |
| 667 | Ok(self.graph) |
| 668 | } |
| 669 | |
| 670 | /// Gets current edge count. |
| 671 | pub fn edge_count(&self) -> usize { |
| 672 | self.graph.edge_count() |
| 673 | } |
| 674 | |
| 675 | /// Reports progress if interval is met. |
| 676 | pub fn report_progress(&self, interval: usize) { |
| 677 | if self.graph.num_edges % interval == 0 { |
| 678 | info!("Mmap graph progress: {} edges written to disk", self.graph.num_edges); |
| 679 | } |
| 680 | } |
| 681 | } |
| 682 |