Oregami
Repositories/oxedyne/fe2o3

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
6use crate::{
7 graph::GraphAccessMethod,
8 mmap_index::{
9 MmapIndex,
10 MmapIndexBuilder,
11 },
12};
13
14use oxedyne_fe2o3_core::{
15 prelude::*,
16 info,
17 warn,
18};
19
20use 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)]
35use std::os::unix::io::AsRawFd;
36
37use 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).
47pub 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
65impl 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.
562pub 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
572impl 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