| 1 | //! Simple test binary for batch_compute_file_indexes. |
| 2 | //! |
| 3 | //! Points at a directory of journal files and indexes them. |
| 4 | |
| 5 | // # 1. Create a mount point |
| 6 | // |
| 7 | // dd if=/dev/zero of=/tmp/slow-disk.img bs=1G count=100 |
| 8 | // LOOP=$(sudo losetup -f --show /tmp/slow-disk.img) |
| 9 | // sudo mkfs.ext4 $LOOP |
| 10 | // sudo mkdir -p /mnt/slow-disk |
| 11 | // sudo mount $LOOP /mnt/slow-disk |
| 12 | // sudo chown $USER:$USER /mnt/slow-disk |
| 13 | // |
| 14 | // # 2. Copy journal files |
| 15 | // cp -r ~/repos/tmp/otel-aws /mnt/slow-disk/ |
| 16 | // |
| 17 | // # 3. Unmount and recreate with delay |
| 18 | // sudo umount /mnt/slow-disk |
| 19 | // SIZE=$(sudo blockdev --getsz $LOOP) |
| 20 | // sudo dmsetup create slow-disk --table "0 $SIZE delay $LOOP 0 50 $LOOP 0 50" |
| 21 | // sudo mount /dev/mapper/slow-disk /mnt/slow-disk |
| 22 | // |
| 23 | // # 4. Now /mnt/slow-disk/otel-aws has your journals on a "slow" disk |
| 24 | // |
| 25 | // # 5. Create slow-io cgroup |
| 26 | // sudo mkdir -p /sys/fs/cgroup/slow-io |
| 27 | // echo "+io" | sudo tee /sys/fs/cgroup/cgroup.controllers |
| 28 | // # Find your device's major:minor (e.g., for nvme0n1) |
| 29 | // cat /sys/block/nvme0n1/dev |
| 30 | // # Let's say it's 259:0, Set a 10MB/s read and write limit |
| 31 | // echo "259:0 rbps=10485760 wbps=10485760" | sudo tee /sys/fs/cgroup/slow-io/io.max |
| 32 | |
| 33 | use journal_engine::{ |
| 34 | Facets, FileIndexCacheBuilder, FileIndexKey, IndexingLimits, QueryTimeRange, |
| 35 | batch_compute_file_indexes, |
| 36 | }; |
| 37 | use journal_index::FieldName; |
| 38 | use journal_registry::{Monitor, Registry}; |
| 39 | use std::env; |
| 40 | use std::path::PathBuf; |
| 41 | use tokio_util::sync::CancellationToken; |
| 42 | |
| 43 | #[allow(unused_imports)] |
| 44 | use tracing::{info, warn}; |
| 45 | |
| 46 | #[tokio::main] |
| 47 | async fn main() -> Result<(), Box<dyn std::error::Error>> { |
| 48 | // Initialize tracing |
| 49 | tracing_subscriber::fmt() |
| 50 | .with_max_level(tracing::Level::DEBUG) |
| 51 | .init(); |
| 52 | |
| 53 | // Get directory from args or use default |
| 54 | let dir = if let Some(arg) = env::args().nth(1) { |
| 55 | PathBuf::from(arg) |
| 56 | } else { |
| 57 | PathBuf::from("/mnt/slow-disk/otel-aws") |
| 58 | }; |
| 59 | |
| 60 | info!("scanning directory: {}", dir.display()); |
| 61 | |
| 62 | // Create registry and scan directory |
| 63 | let (monitor, _event_receiver) = Monitor::new()?; |
| 64 | let registry = Registry::new(monitor); |
| 65 | |
| 66 | registry.watch_directory(dir.to_str().unwrap())?; |
| 67 | |
| 68 | // Find all files |
| 69 | let files = registry.find_files_in_range( |
| 70 | journal_common::Seconds(0), |
| 71 | journal_common::Seconds(u32::MAX), |
| 72 | )?; |
| 73 | |
| 74 | info!("found {} journal files", files.len()); |
| 75 | if files.is_empty() { |
| 76 | return Ok(()); |
| 77 | } |
| 78 | // files.truncate(1); |
| 79 | |
| 80 | // Create file index cache |
| 81 | let cache = FileIndexCacheBuilder::new() |
| 82 | // .with_cache_path("/mnt/slow-disk/foyer-cache") |
| 83 | .with_cache_path("/tmp/foyer-cache") |
| 84 | .with_memory_capacity(1000) |
| 85 | .with_disk_capacity(2048 * 1024 * 1024) |
| 86 | .with_block_size(4 * 1024 * 1024) |
| 87 | .build() |
| 88 | .await?; |
| 89 | |
| 90 | info!("created file index cache"); |
| 91 | |
| 92 | // Configure indexing parameters (modify these as needed) |
| 93 | let facets = Facets::new(&["log.severity_number".to_string()]); |
| 94 | let source_timestamp_field = FieldName::new("_SOURCE_REALTIME_TIMESTAMP").unwrap(); |
| 95 | |
| 96 | let keys: Vec<FileIndexKey> = files |
| 97 | .iter() |
| 98 | .map(|file_info| { |
| 99 | FileIndexKey::new( |
| 100 | &file_info.file, |
| 101 | &facets, |
| 102 | Some(source_timestamp_field.clone()), |
| 103 | ) |
| 104 | }) |
| 105 | .collect(); |
| 106 | |
| 107 | // Create a time range for indexing (24 hours) |
| 108 | let now = std::time::SystemTime::now() |
| 109 | .duration_since(std::time::UNIX_EPOCH)? |
| 110 | .as_secs() as u32; |
| 111 | let time_range = QueryTimeRange::new(now - 86400, now)?; |
| 112 | let cancellation = CancellationToken::new(); |
| 113 | |
| 114 | info!( |
| 115 | "computing {} file indexes, bucket duration: {}s", |
| 116 | keys.len(), |
| 117 | time_range.bucket_duration() |
| 118 | ); |
| 119 | |
| 120 | // Run batch indexing |
| 121 | let start = std::time::Instant::now(); |
| 122 | let responses = batch_compute_file_indexes( |
| 123 | &cache, |
| 124 | ®istry, |
| 125 | keys, |
| 126 | &time_range, |
| 127 | cancellation, |
| 128 | IndexingLimits::default(), |
| 129 | None, |
| 130 | ) |
| 131 | .await?; |
| 132 | |
| 133 | let elapsed = start.elapsed(); |
| 134 | |
| 135 | info!("responses={}, duration={:?}", responses.len(), elapsed); |
| 136 | |
| 137 | // Close the cache to flush and shut down I/O tasks gracefully |
| 138 | cache.close().await?; |
| 139 | |
| 140 | Ok(()) |
| 141 | } |