use std::cell::RefCell;
use std::collections::{BTreeMap, BTreeSet};
use std::ffi::{OsStr, OsString};
use std::fmt;
use std::io::Write;
use std::mem::{offset_of, size_of};
use std::os::windows::ffi::{OsStrExt, OsStringExt};
use std::os::windows::fs::OpenOptionsExt;
use std::os::windows::io::AsRawHandle;
use std::path::{Component, Path, PathBuf, Prefix};
use std::ptr;
use std::sync::{Arc, Condvar, Mutex, OnceLock};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use rayon::{ThreadPool, prelude::*};
use rebecca_ntfs::{
MftIndex, MftIndexEntry, MftRecordBatch, MftRecordReader, NtfsDirectoryEntry,
NtfsFileReference, NtfsParsedRecord, NtfsRecordSet, NtfsStreamGeometry, NtfsStreamSource,
ParseCaveat, PhysicalMetrics, PhysicalMetricsAccumulator, SubtreeSummary,
resolve_record_with_stream_source,
};
use serde::{Deserialize, Serialize};
use windows::Win32::Foundation::{
CloseHandle, ERROR_ACCESS_DENIED, ERROR_HANDLE_EOF, ERROR_INVALID_PARAMETER, ERROR_MORE_DATA,
HANDLE, WIN32_ERROR,
};
use windows::Win32::Storage::FileSystem::{
BY_HANDLE_FILE_INFORMATION, CreateFileW, FILE_BEGIN, FILE_FLAG_BACKUP_SEMANTICS,
FILE_FLAG_OPEN_REPARSE_POINT, FILE_FLAG_SEQUENTIAL_SCAN, FILE_FLAGS_AND_ATTRIBUTES,
FILE_READ_ATTRIBUTES, FILE_SHARE_DELETE, FILE_SHARE_MODE, FILE_SHARE_READ, FILE_SHARE_WRITE,
GetDriveTypeW, GetFileInformationByHandle, GetVolumeInformationW, OPEN_EXISTING, ReadFile,
SYNCHRONIZE, SetFilePointerEx,
};
use windows::Win32::System::IO::DeviceIoControl;
use windows::Win32::System::Ioctl::{
FSCTL_GET_NTFS_FILE_RECORD, FSCTL_GET_NTFS_VOLUME_DATA, FSCTL_GET_RETRIEVAL_POINTERS,
FSCTL_QUERY_USN_JOURNAL, FSCTL_READ_USN_JOURNAL, NTFS_FILE_RECORD_INPUT_BUFFER,
NTFS_FILE_RECORD_OUTPUT_BUFFER, NTFS_VOLUME_DATA_BUFFER, READ_USN_JOURNAL_DATA_V0,
RETRIEVAL_POINTERS_BUFFER, RETRIEVAL_POINTERS_BUFFER_0, STARTING_VCN_INPUT_BUFFER,
USN_JOURNAL_DATA_V0, USN_RECORD_V2,
};
use windows::core::{Error as WindowsError, HRESULT, PCWSTR};
use crate::disk_map::{
DiskMapBackendOptions, DiskMapBackendReport, DiskMapEntry, DiskMapEntryKind,
DiskMapFileIdentity, DiskMapGroupCollector, DiskMapMetadataSemantics, DiskMapMetrics,
DiskMapTopEntries,
};
use crate::error::{RebeccaError, Result, ScanFailure, ScanFailurePhase};
use crate::parallelism::{bounded_parallelism_budget, run_scoped_parallel_work};
use crate::plan::{EstimateProvenance, EstimateSource};
use crate::progress::{InspectProgressEvent, InspectProgressResult};
use crate::safety::is_reparse_like;
use crate::scan_cache::{ScanCacheUsnCheckpoint, ScanCacheUsnJournalState};
use super::backend::{
MeasuredScan, ScanBackend, ScanBackendEvidence, ScanBackendKind, ScanEstimateConfidence,
ScanRequest,
};
use super::progress::{ScanProgressEvent, check_not_cancelled};
use super::{ScanCancellationToken, ScanReport};
const EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL: &str = "windows-ntfs-mft-experimental";
const NTFS_FILE_SYSTEM_NAME: &str = "NTFS";
const DRIVE_FIXED: u32 = 3;
const FILE_REFERENCE_LOW_MASK: u64 = 0x0000_FFFF_FFFF_FFFF;
const SEQUENTIAL_MFT_SOURCE_LABEL: &str = "sequential";
const FSCTL_RECORD_SOURCE_LABEL: &str = "fsctl-record";
const TARGETED_MFT_SOURCE_LABEL: &str = "targeted-fsctl";
const PERSISTENT_MFT_SOURCE_LABEL: &str = "persistent-cache";
const SEQUENTIAL_MFT_CHUNK_BYTES: usize = 8 * 1024 * 1024;
const SEQUENTIAL_MFT_PARSE_WINDOW_CHUNKS: usize = 8;
const MFT_MIRROR_SYSTEM_RECORDS: u64 = 4;
const TARGETED_MFT_MAX_RECORDS: usize = 1_000_000;
const TARGETED_MFT_MAX_DEPTH: usize = 512;
const NTFS_FILE_RECORD_OUTPUT_HEADER_BYTES: usize = size_of::<i64>() + size_of::<u32>();
const MAX_RETRIEVAL_POINTER_BUFFER_BYTES: usize = 16 * 1024 * 1024;
const USN_JOURNAL_READ_BUFFER_BYTES: usize = 1024 * 1024;
const MAX_NTFS_VOLUME_INDEX_USN_REPLAY_RECORDS: usize = 100_000;
const MAX_NTFS_VOLUME_INDEX_USN_REPLAY_ATTEMPTS: usize = 3;
const MAX_MFT_PARSE_ERROR_CAVEAT_SAMPLES: usize = 8;
const MAX_MFT_ESTIMATE_CAVEAT_SAMPLES_PER_CODE: usize = 8;
const MFT_CAVEAT_SUMMARY_CODE: &str = "mft-caveat-summary";
const MFT_MIRROR_READ_FAILED_CAVEAT_CODE: &str = "mft-mirror-read-failed";
const LIVE_NTFS_MFT_INDEX_TIMEOUT_ENV: &str = "REBECCA_NTFS_MFT_INDEX_TIMEOUT_SECONDS";
const LIVE_NTFS_MFT_INDEX_TIMINGS_ENV: &str = "REBECCA_NTFS_MFT_INDEX_TIMINGS";
const LIVE_NTFS_MFT_FULL_INDEX_FALLBACK_ENV: &str = "REBECCA_NTFS_MFT_FULL_INDEX_FALLBACK";
const DEFAULT_LIVE_NTFS_MFT_INDEX_TIMEOUT: Duration = Duration::from_secs(20);
const MFT_BUILD_TIMING_CAVEAT_CODE: &str = "mft-index-build-timing";
const MFT_INDEX_ALLOCATION_BUDGET_EXHAUSTED_CAVEAT_CODE: &str =
"mft-index-allocation-budget-exhausted";
const MFT_PERSISTENT_CACHE_MISS_CAVEAT_CODE: &str = "mft-persistent-cache-miss";
const MFT_PERSISTENT_CACHE_WRITE_SKIPPED_CAVEAT_CODE: &str = "mft-persistent-cache-write-skipped";
const NTFS_VOLUME_INDEX_MANIFEST_VERSION: u32 = 1;
const NTFS_VOLUME_INDEX_PAYLOAD_VERSION: u32 = 1;
const NTFS_VOLUME_INDEX_CACHE_DIR: &str = "ntfs-volume-index";
static MFT_PARSE_THREAD_POOL: OnceLock<ThreadPool> = OnceLock::new();
#[derive(Debug, Clone, Copy)]
struct NtfsMftBuildBudget {
started_at: Instant,
timeout: Option<Duration>,
}
impl NtfsMftBuildBudget {
fn new(timeout: Option<Duration>) -> Self {
Self {
started_at: Instant::now(),
timeout,
}
}
}
#[derive(Debug, Default)]
struct NtfsMftBuildMonitorState {
active_stage: Option<(NtfsMftBuildStage, Instant)>,
timings: BTreeMap<NtfsMftBuildStage, Duration>,
metrics: BTreeMap<NtfsMftBuildMetric, u64>,
observer_error: Option<RebeccaError>,
}
#[derive(Debug, Clone, Copy)]
enum NtfsMftBuildMonitorEvent {
StageStarted { stage: &'static str },
StageFinished { stage: &'static str },
Metric { metric: &'static str, value: u64 },
}
type NtfsMftBuildObserver<'a> = dyn FnMut(NtfsMftBuildMonitorEvent) -> Result<()> + 'a;
struct NtfsMftBuildMonitor<'a> {
budget: NtfsMftBuildBudget,
state: RefCell<NtfsMftBuildMonitorState>,
emit_timing_caveat: bool,
observer: Option<RefCell<&'a mut NtfsMftBuildObserver<'a>>>,
}
impl<'a> NtfsMftBuildMonitor<'a> {
fn from_environment() -> Self {
Self::new(
live_ntfs_mft_index_timeout(),
live_ntfs_mft_index_timings_enabled(),
)
}
fn from_environment_with_observer(observer: &'a mut NtfsMftBuildObserver<'a>) -> Self {
Self::new_with_observer(
live_ntfs_mft_index_timeout(),
live_ntfs_mft_index_timings_enabled(),
observer,
)
}
fn new(timeout: Option<Duration>, emit_timing_caveat: bool) -> Self {
Self {
budget: NtfsMftBuildBudget::new(timeout),
state: RefCell::new(NtfsMftBuildMonitorState::default()),
emit_timing_caveat,
observer: None,
}
}
fn new_with_observer(
timeout: Option<Duration>,
emit_timing_caveat: bool,
observer: &'a mut NtfsMftBuildObserver<'a>,
) -> Self {
Self {
budget: NtfsMftBuildBudget::new(timeout),
state: RefCell::new(NtfsMftBuildMonitorState::default()),
emit_timing_caveat,
observer: Some(RefCell::new(observer)),
}
}
#[cfg(test)]
fn measure<R>(
&self,
stage: NtfsMftBuildStage,
operation: impl FnOnce() -> Result<R>,
) -> Result<R> {
self.measure_inner(stage, operation, || Ok(()))
}
fn measure_checked<R>(
&self,
stage: NtfsMftBuildStage,
cancellation: &ScanCancellationToken,
operation: impl FnOnce() -> Result<R>,
) -> Result<R> {
self.check(cancellation)?;
self.measure_inner(stage, operation, || self.check(cancellation))
}
fn measure_allow_budget_overrun_after_success<R>(
&self,
stage: NtfsMftBuildStage,
cancellation: &ScanCancellationToken,
operation: impl FnOnce() -> Result<R>,
) -> Result<R> {
self.check(cancellation)?;
self.measure_inner(stage, operation, || check_not_cancelled(cancellation))
}
fn measure_cancellation_only<R>(
&self,
stage: NtfsMftBuildStage,
cancellation: &ScanCancellationToken,
operation: impl FnOnce() -> Result<R>,
) -> Result<R> {
check_not_cancelled(cancellation)?;
self.measure_inner(stage, operation, || check_not_cancelled(cancellation))
}
fn measure_inner<R>(
&self,
stage: NtfsMftBuildStage,
operation: impl FnOnce() -> Result<R>,
after_success: impl FnOnce() -> Result<()>,
) -> Result<R> {
let started_at = Instant::now();
let previous_stage = {
let mut state = self.state.borrow_mut();
state.active_stage.replace((stage, started_at))
};
if let Err(err) = self.emit_event(NtfsMftBuildMonitorEvent::StageStarted {
stage: stage.label(),
}) {
let mut state = self.state.borrow_mut();
state.active_stage = previous_stage;
return Err(err);
}
let result = operation();
let elapsed = started_at.elapsed();
let after_success = if result.is_ok() {
after_success()
} else {
Ok(())
};
{
let mut state = self.state.borrow_mut();
let total = state.timings.entry(stage).or_default();
*total = total.saturating_add(elapsed);
state.active_stage = previous_stage;
}
let stage_finished = self.emit_event(NtfsMftBuildMonitorEvent::StageFinished {
stage: stage.label(),
});
match result {
Ok(value) => {
after_success?;
stage_finished?;
self.check_observer_error()?;
Ok(value)
}
Err(err) => Err(err),
}
}
fn check(&self, cancellation: &ScanCancellationToken) -> Result<()> {
check_not_cancelled(cancellation)?;
self.check_observer_error()?;
if !self.is_timed_out() {
return Ok(());
}
let Some(timeout) = self.budget.timeout else {
return Ok(());
};
let state = self.state.borrow();
let stage = state
.active_stage
.map(|(stage, started_at)| {
format!(
" while {}; stage_elapsed={}ms",
stage.label(),
started_at.elapsed().as_millis()
)
})
.unwrap_or_default();
let timings = format_build_summary(&state.timings, &state.metrics)
.map(|summary| format!("; {summary}"))
.unwrap_or_default();
Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} live volume index build timed out after {}s{stage}{timings}; tune {LIVE_NTFS_MFT_INDEX_TIMEOUT_ENV} to increase the budget or set it to 0 to disable this guard",
timeout.as_secs()
)))
}
fn is_timed_out(&self) -> bool {
self.budget
.timeout
.is_some_and(|timeout| self.budget.started_at.elapsed() >= timeout)
}
fn timing_caveat(&self) -> Option<ParseCaveat> {
if !self.emit_timing_caveat {
return None;
}
self.build_summary().map(|summary| {
ParseCaveat::new(
MFT_BUILD_TIMING_CAVEAT_CODE,
format!("live NTFS/MFT index build timings: {summary}"),
)
})
}
fn build_summary(&self) -> Option<String> {
let state = self.state.borrow();
format_build_summary(&state.timings, &state.metrics)
}
fn evidence(&self) -> ScanBackendEvidence {
let state = self.state.borrow();
ScanBackendEvidence {
timings_ms: state
.timings
.iter()
.map(|(stage, duration)| {
(
stage.label().to_string(),
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX),
)
})
.collect(),
counters: state
.metrics
.iter()
.filter(|(_, value)| **value > 0)
.map(|(metric, value)| (metric.label().to_string(), *value))
.collect(),
cache_events: Vec::new(),
}
}
fn add_metric(&self, metric: NtfsMftBuildMetric, value: u64) {
if value == 0 {
return;
}
let value = {
let mut state = self.state.borrow_mut();
let total = state.metrics.entry(metric).or_default();
*total = total.saturating_add(value);
*total
};
if let Err(err) = self.emit_event(NtfsMftBuildMonitorEvent::Metric {
metric: metric.label(),
value,
}) {
self.state.borrow_mut().observer_error = Some(err);
}
}
fn add_metric_usize(&self, metric: NtfsMftBuildMetric, value: usize) {
self.add_metric(metric, u64::try_from(value).unwrap_or(u64::MAX));
}
fn emit_event(&self, event: NtfsMftBuildMonitorEvent) -> Result<()> {
self.check_observer_error()?;
let Some(observer) = &self.observer else {
return Ok(());
};
observer.borrow_mut()(event)
}
fn check_observer_error(&self) -> Result<()> {
if let Some(err) = self.state.borrow_mut().observer_error.take() {
return Err(err);
}
Ok(())
}
#[cfg(test)]
fn expired_for_test(timeout: Duration) -> Self {
Self::elapsed_for_test(timeout, timeout + Duration::from_secs(1))
}
#[cfg(test)]
fn near_timeout_for_test(timeout: Duration, remaining: Duration) -> Self {
Self::elapsed_for_test(timeout, timeout.saturating_sub(remaining))
}
#[cfg(test)]
fn elapsed_for_test(timeout: Duration, elapsed: Duration) -> Self {
let mut monitor = Self::new(Some(timeout), false);
let now = Instant::now();
monitor.budget.started_at = now.checked_sub(elapsed).unwrap_or(now);
monitor
}
}
fn ntfs_mft_build_monitor<'a>(
observer: Option<&'a mut NtfsMftBuildObserver<'a>>,
) -> NtfsMftBuildMonitor<'a> {
match observer {
Some(observer) => NtfsMftBuildMonitor::from_environment_with_observer(observer),
None => NtfsMftBuildMonitor::from_environment(),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
enum NtfsMftBuildStage {
OpenVolume,
ReadVolumeData,
SequentialOpenMftData,
SequentialReadRetrievalPointers,
SequentialReadMftBytes,
SequentialReadMftMirror,
SequentialParseRecords,
FsctlReadParseRecords,
TargetedReadRecord,
TargetedResolveRecord,
TargetedTraverseSubtree,
ResolveIndexAllocations,
BuildMftIndex,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
enum NtfsMftBuildMetric {
ParsedRecords,
SequentialMftReadBytes,
SequentialMftReadChunks,
MftMirrorReadBytes,
MftMirrorReadChunks,
FsctlRecordAttempts,
FsctlRecordSuccesses,
TargetedRecordAttempts,
TargetedRecordSuccesses,
StreamReadBytes,
StreamReads,
}
impl NtfsMftBuildMetric {
fn label(self) -> &'static str {
match self {
Self::ParsedRecords => "parsed-records",
Self::SequentialMftReadBytes => "sequential-mft-read-bytes",
Self::SequentialMftReadChunks => "sequential-mft-read-chunks",
Self::MftMirrorReadBytes => "mft-mirror-read-bytes",
Self::MftMirrorReadChunks => "mft-mirror-read-chunks",
Self::FsctlRecordAttempts => "fsctl-record-attempts",
Self::FsctlRecordSuccesses => "fsctl-record-successes",
Self::TargetedRecordAttempts => "targeted-record-attempts",
Self::TargetedRecordSuccesses => "targeted-record-successes",
Self::StreamReadBytes => "stream-read-bytes",
Self::StreamReads => "stream-reads",
}
}
}
impl NtfsMftBuildStage {
fn label(self) -> &'static str {
match self {
Self::OpenVolume => "open-volume",
Self::ReadVolumeData => "read-volume-data",
Self::SequentialOpenMftData => "sequential-open-mft-data",
Self::SequentialReadRetrievalPointers => "sequential-read-retrieval-pointers",
Self::SequentialReadMftBytes => "sequential-read-mft-bytes",
Self::SequentialReadMftMirror => "sequential-read-mft-mirror",
Self::SequentialParseRecords => "sequential-parse-records",
Self::FsctlReadParseRecords => "fsctl-read-parse-records",
Self::TargetedReadRecord => "targeted-read-record",
Self::TargetedResolveRecord => "targeted-resolve-record",
Self::TargetedTraverseSubtree => "targeted-traverse-subtree",
Self::ResolveIndexAllocations => "resolve-index-allocations",
Self::BuildMftIndex => "build-mft-index",
}
}
}
fn format_timing_summary(timings: &BTreeMap<NtfsMftBuildStage, Duration>) -> Option<String> {
if timings.is_empty() {
return None;
}
Some(
timings
.iter()
.map(|(stage, duration)| format!("{}={}ms", stage.label(), duration.as_millis()))
.collect::<Vec<_>>()
.join(", "),
)
}
fn format_metric_summary(metrics: &BTreeMap<NtfsMftBuildMetric, u64>) -> Option<String> {
let values: Vec<_> = metrics
.iter()
.filter(|(_, value)| **value > 0)
.map(|(metric, value)| format!("{}={value}", metric.label()))
.collect();
if values.is_empty() {
return None;
}
Some(values.join(", "))
}
fn format_build_summary(
timings: &BTreeMap<NtfsMftBuildStage, Duration>,
metrics: &BTreeMap<NtfsMftBuildMetric, u64>,
) -> Option<String> {
let mut parts = Vec::new();
if let Some(timings) = format_timing_summary(timings) {
parts.push(format!("completed_timings={timings}"));
}
if let Some(metrics) = format_metric_summary(metrics) {
parts.push(format!("metrics={metrics}"));
}
if parts.is_empty() {
return None;
}
Some(parts.join("; "))
}
fn live_ntfs_mft_index_timeout() -> Option<Duration> {
let Some(raw) = std::env::var_os(LIVE_NTFS_MFT_INDEX_TIMEOUT_ENV) else {
return Some(DEFAULT_LIVE_NTFS_MFT_INDEX_TIMEOUT);
};
let raw = raw.to_string_lossy();
let trimmed = raw.trim();
if trimmed.is_empty() {
return Some(DEFAULT_LIVE_NTFS_MFT_INDEX_TIMEOUT);
}
match trimmed.parse::<u64>() {
Ok(0) => None,
Ok(seconds) => Some(Duration::from_secs(seconds)),
Err(_) => Some(DEFAULT_LIVE_NTFS_MFT_INDEX_TIMEOUT),
}
}
fn live_ntfs_mft_index_timings_enabled() -> bool {
std::env::var_os(LIVE_NTFS_MFT_INDEX_TIMINGS_ENV).is_some_and(|raw| {
let raw = raw.to_string_lossy();
let trimmed = raw.trim();
!trimmed.is_empty() && trimmed != "0"
})
}
fn live_ntfs_mft_full_index_fallback_enabled() -> bool {
std::env::var_os(LIVE_NTFS_MFT_FULL_INDEX_FALLBACK_ENV).is_some_and(|raw| {
let raw = raw.to_string_lossy();
let trimmed = raw.trim();
!trimmed.is_empty() && trimmed != "0"
})
}
fn check_mft_build_progress(
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<()> {
monitor.check(cancellation)
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
struct NtfsVolumeIndexCacheKey {
device_path: String,
volume_serial: u64,
}
impl NtfsVolumeIndexCacheKey {
fn new(device_path: impl Into<String>, volume_serial: u64) -> Self {
Self {
device_path: device_path.into(),
volume_serial,
}
}
}
impl fmt::Display for NtfsVolumeIndexCacheKey {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "{}#{}", self.device_path, self.volume_serial)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct NtfsVolumeIndexFingerprint {
key: NtfsVolumeIndexCacheKey,
record_size: u64,
sector_size: u64,
bytes_per_cluster: u64,
mft_start_lcn: u64,
mft_mirror_start_lcn: u64,
mft_valid_data_length: u64,
}
impl NtfsVolumeIndexFingerprint {
fn from_volume_data(
capabilities: &NtfsVolumeCapabilities,
volume_data: &NTFS_VOLUME_DATA_BUFFER,
geometry: NtfsRecordGeometry,
) -> Self {
Self {
key: capabilities.cache_key(),
record_size: geometry.record_size as u64,
sector_size: geometry.sector_size as u64,
bytes_per_cluster: geometry.bytes_per_cluster,
mft_start_lcn: u64::try_from(volume_data.MftStartLcn).unwrap_or(0),
mft_mirror_start_lcn: u64::try_from(volume_data.Mft2StartLcn).unwrap_or(0),
mft_valid_data_length: u64::try_from(volume_data.MftValidDataLength).unwrap_or(0),
}
}
fn is_reusable_for(&self, capabilities: &NtfsVolumeCapabilities) -> bool {
self.key.device_path == capabilities.device_path
&& self.key.volume_serial == capabilities.volume_serial
&& self.persistent_generation() != 0
}
fn persistent_generation(&self) -> u64 {
const SCHEMA_VERSION: u64 = 1;
let mut hash = FNV_OFFSET_BASIS;
update_u64(&mut hash, SCHEMA_VERSION);
update_bytes(&mut hash, self.key.device_path.as_bytes());
update_u64(&mut hash, self.key.volume_serial);
update_u64(&mut hash, self.record_size);
update_u64(&mut hash, self.sector_size);
update_u64(&mut hash, self.bytes_per_cluster);
update_u64(&mut hash, self.mft_start_lcn);
update_u64(&mut hash, self.mft_mirror_start_lcn);
update_u64(&mut hash, self.mft_valid_data_length);
hash
}
}
const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
fn update_bytes(hash: &mut u64, bytes: &[u8]) {
for byte in bytes {
*hash ^= u64::from(*byte);
*hash = hash.wrapping_mul(FNV_PRIME);
}
}
fn update_u64(hash: &mut u64, value: u64) {
update_bytes(hash, &value.to_le_bytes());
}
fn cache_checksum(raw: &[u8]) -> u64 {
let mut hash = FNV_OFFSET_BASIS;
update_u64(&mut hash, NTFS_VOLUME_INDEX_PAYLOAD_VERSION as u64);
update_bytes(&mut hash, raw);
hash
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct NtfsVolumeIndexPayloadRef {
version: u32,
file_name: String,
byte_len: u64,
checksum: u64,
}
impl NtfsVolumeIndexPayloadRef {
fn new(file_name: String, raw: &[u8]) -> Self {
Self {
version: NTFS_VOLUME_INDEX_PAYLOAD_VERSION,
file_name,
byte_len: raw.len() as u64,
checksum: cache_checksum(raw),
}
}
fn validation_miss(&self, raw: &[u8]) -> Option<NtfsVolumeIndexPayloadMiss> {
if self.version != NTFS_VOLUME_INDEX_PAYLOAD_VERSION {
return Some(NtfsVolumeIndexPayloadMiss::UnsupportedVersion);
}
if self.byte_len != raw.len() as u64 || self.checksum != cache_checksum(raw) {
return Some(NtfsVolumeIndexPayloadMiss::ChecksumMismatch);
}
None
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct NtfsVolumeIndexPayload {
version: u32,
generation: u64,
fingerprint: NtfsVolumeIndexFingerprint,
source_label: String,
caveats: Vec<ParseCaveat>,
mft_index: MftIndex,
}
impl NtfsVolumeIndexPayload {
fn new(
fingerprint: NtfsVolumeIndexFingerprint,
source_label: &'static str,
caveats: Vec<ParseCaveat>,
mft_index: MftIndex,
) -> Self {
Self {
version: NTFS_VOLUME_INDEX_PAYLOAD_VERSION,
generation: fingerprint.persistent_generation(),
fingerprint,
source_label: source_label.to_string(),
caveats,
mft_index,
}
}
fn validation_miss(
&self,
fingerprint: &NtfsVolumeIndexFingerprint,
) -> Option<NtfsVolumeIndexPayloadMiss> {
if self.version != NTFS_VOLUME_INDEX_PAYLOAD_VERSION {
return Some(NtfsVolumeIndexPayloadMiss::UnsupportedVersion);
}
if self.generation != fingerprint.persistent_generation()
|| self.fingerprint != *fingerprint
{
return Some(NtfsVolumeIndexPayloadMiss::Stale);
}
None
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct NtfsVolumeIndexCacheManifest {
version: u32,
generation: u64,
fingerprint: NtfsVolumeIndexFingerprint,
source_label: String,
created_at_unix_seconds: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
payload: Option<NtfsVolumeIndexPayloadRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
usn_checkpoint: Option<ScanCacheUsnCheckpoint>,
}
impl NtfsVolumeIndexCacheManifest {
fn new(
fingerprint: NtfsVolumeIndexFingerprint,
source_label: &'static str,
usn_checkpoint: Option<ScanCacheUsnCheckpoint>,
) -> Self {
Self {
version: NTFS_VOLUME_INDEX_MANIFEST_VERSION,
generation: fingerprint.persistent_generation(),
fingerprint,
source_label: source_label.to_string(),
created_at_unix_seconds: unix_now(),
payload: None,
usn_checkpoint,
}
}
fn with_payload(mut self, payload: NtfsVolumeIndexPayloadRef) -> Self {
self.payload = Some(payload);
self
}
fn validation_miss(
&self,
fingerprint: &NtfsVolumeIndexFingerprint,
) -> Option<NtfsVolumeIndexManifestMiss> {
if self.version != NTFS_VOLUME_INDEX_MANIFEST_VERSION {
return Some(NtfsVolumeIndexManifestMiss::UnsupportedVersion);
}
if self.generation != fingerprint.persistent_generation()
|| self.fingerprint != *fingerprint
{
return Some(NtfsVolumeIndexManifestMiss::Stale);
}
None
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum NtfsVolumeIndexManifestMiss {
Missing,
Unreadable,
Corrupted,
UnsupportedVersion,
Stale,
}
impl NtfsVolumeIndexManifestMiss {
fn should_prune(self) -> bool {
matches!(self, Self::Unreadable | Self::Corrupted | Self::Stale)
}
const fn label(self) -> &'static str {
match self {
Self::Missing => "manifest-missing",
Self::Unreadable => "manifest-unreadable",
Self::Corrupted => "manifest-corrupted",
Self::UnsupportedVersion => "manifest-unsupported-version",
Self::Stale => "manifest-stale",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum NtfsVolumeIndexManifestLookup {
Hit(NtfsVolumeIndexCacheManifest),
Miss(NtfsVolumeIndexManifestMiss),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum NtfsVolumeIndexPayloadMiss {
Manifest(NtfsVolumeIndexManifestMiss),
ManifestWithoutPayload,
Missing,
Unreadable,
Corrupted,
UnsupportedVersion,
Stale,
ChecksumMismatch,
UsnCheckpointMissing,
UsnJournalChanged,
UsnRangeUnavailable,
UsnReplayUnstable,
UsnAncestryUnavailable,
UsnTargetChanged,
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum NtfsVolumeIndexPayloadLookup {
Hit {
manifest: Box<NtfsVolumeIndexCacheManifest>,
payload: Box<NtfsVolumeIndexPayload>,
},
Miss(NtfsVolumeIndexPayloadMiss),
}
impl NtfsVolumeIndexPayloadMiss {
fn should_prune_pair(self) -> bool {
match self {
Self::Manifest(reason) => reason.should_prune(),
Self::ManifestWithoutPayload | Self::UnsupportedVersion => false,
Self::Missing
| Self::Unreadable
| Self::Corrupted
| Self::Stale
| Self::ChecksumMismatch
| Self::UsnJournalChanged
| Self::UsnRangeUnavailable => true,
Self::UsnCheckpointMissing
| Self::UsnReplayUnstable
| Self::UsnAncestryUnavailable
| Self::UsnTargetChanged => false,
}
}
fn label(&self) -> String {
match self {
Self::Manifest(reason) => reason.label().to_string(),
Self::ManifestWithoutPayload => "manifest-without-payload".to_string(),
Self::Missing => "payload-missing".to_string(),
Self::Unreadable => "payload-unreadable".to_string(),
Self::Corrupted => "payload-corrupted".to_string(),
Self::UnsupportedVersion => "payload-unsupported-version".to_string(),
Self::Stale => "payload-stale".to_string(),
Self::ChecksumMismatch => "payload-checksum-mismatch".to_string(),
Self::UsnCheckpointMissing => "usn-checkpoint-missing".to_string(),
Self::UsnJournalChanged => "usn-journal-changed".to_string(),
Self::UsnRangeUnavailable => "usn-range-unavailable".to_string(),
Self::UsnReplayUnstable => "usn-replay-unstable".to_string(),
Self::UsnAncestryUnavailable => "usn-ancestry-unavailable".to_string(),
Self::UsnTargetChanged => "usn-target-changed".to_string(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) struct NtfsVolumeIndexManifestStore {
root_dir: PathBuf,
}
impl NtfsVolumeIndexManifestStore {
pub(super) fn new(cache_root: impl Into<PathBuf>) -> Self {
Self {
root_dir: cache_root.into().join(NTFS_VOLUME_INDEX_CACHE_DIR),
}
}
fn cache_file_for(&self, fingerprint: &NtfsVolumeIndexFingerprint) -> PathBuf {
self.root_dir
.join(format!("{:016x}.json", fingerprint.persistent_generation()))
}
fn payload_file_name_for(&self, fingerprint: &NtfsVolumeIndexFingerprint) -> String {
format!("{:016x}.index.json", fingerprint.persistent_generation())
}
fn payload_file_for(&self, fingerprint: &NtfsVolumeIndexFingerprint) -> PathBuf {
self.root_dir.join(self.payload_file_name_for(fingerprint))
}
fn load(&self, fingerprint: &NtfsVolumeIndexFingerprint) -> NtfsVolumeIndexManifestLookup {
let cache_file = self.cache_file_for(fingerprint);
let raw = match std::fs::read_to_string(&cache_file) {
Ok(raw) => raw,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
return NtfsVolumeIndexManifestLookup::Miss(NtfsVolumeIndexManifestMiss::Missing);
}
Err(_) => {
prune_manifest_file(&cache_file);
return NtfsVolumeIndexManifestLookup::Miss(
NtfsVolumeIndexManifestMiss::Unreadable,
);
}
};
let manifest: NtfsVolumeIndexCacheManifest = match serde_json::from_str(&raw) {
Ok(manifest) => manifest,
Err(_) => {
prune_manifest_file(&cache_file);
return NtfsVolumeIndexManifestLookup::Miss(NtfsVolumeIndexManifestMiss::Corrupted);
}
};
if let Some(reason) = manifest.validation_miss(fingerprint) {
if reason.should_prune() {
prune_manifest_file(&cache_file);
}
return NtfsVolumeIndexManifestLookup::Miss(reason);
}
NtfsVolumeIndexManifestLookup::Hit(manifest)
}
fn load_index_payload(
&self,
fingerprint: &NtfsVolumeIndexFingerprint,
) -> NtfsVolumeIndexPayloadLookup {
let manifest = match self.load(fingerprint) {
NtfsVolumeIndexManifestLookup::Hit(manifest) => manifest,
NtfsVolumeIndexManifestLookup::Miss(reason) => {
return NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Manifest(
reason,
));
}
};
let Some(payload_ref) = manifest.payload.as_ref() else {
return NtfsVolumeIndexPayloadLookup::Miss(
NtfsVolumeIndexPayloadMiss::ManifestWithoutPayload,
);
};
let manifest_file = self.cache_file_for(fingerprint);
let payload_file = self.payload_file_for(fingerprint);
let expected_file_name = self.payload_file_name_for(fingerprint);
if payload_ref.file_name != expected_file_name {
prune_manifest_file(&manifest_file);
return NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Stale);
}
let raw = match std::fs::read(&payload_file) {
Ok(raw) => raw,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
prune_manifest_file(&manifest_file);
return NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Missing);
}
Err(_) => {
prune_manifest_pair(&manifest_file, &payload_file);
return NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Unreadable);
}
};
if let Some(reason) = payload_ref.validation_miss(&raw) {
if reason.should_prune_pair() {
prune_manifest_pair(&manifest_file, &payload_file);
}
return NtfsVolumeIndexPayloadLookup::Miss(reason);
}
let payload: NtfsVolumeIndexPayload = match serde_json::from_slice(&raw) {
Ok(payload) => payload,
Err(_) => {
prune_manifest_pair(&manifest_file, &payload_file);
return NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Corrupted);
}
};
if let Some(reason) = payload.validation_miss(fingerprint) {
if reason.should_prune_pair() {
prune_manifest_pair(&manifest_file, &payload_file);
}
return NtfsVolumeIndexPayloadLookup::Miss(reason);
}
NtfsVolumeIndexPayloadLookup::Hit {
manifest: Box::new(manifest),
payload: Box::new(payload),
}
}
fn load_replayable_index_payload(
&self,
fingerprint: &NtfsVolumeIndexFingerprint,
journal_state: &ScanCacheUsnJournalState,
) -> NtfsVolumeIndexPayloadLookup {
let lookup = self.load_index_payload(fingerprint);
let NtfsVolumeIndexPayloadLookup::Hit { manifest, payload } = lookup else {
return lookup;
};
if let Some(reason) =
validate_ntfs_volume_index_replay_range(manifest.usn_checkpoint.as_ref(), journal_state)
{
if reason.should_prune_pair() {
prune_manifest_pair(
&self.cache_file_for(fingerprint),
&self.payload_file_for(fingerprint),
);
}
return NtfsVolumeIndexPayloadLookup::Miss(reason);
}
NtfsVolumeIndexPayloadLookup::Hit { manifest, payload }
}
fn store(&self, manifest: &NtfsVolumeIndexCacheManifest) -> Result<()> {
let cache_file = self.cache_file_for(&manifest.fingerprint);
let raw = serde_json::to_vec(manifest)?;
write_ntfs_volume_index_file(&cache_file, &raw, "manifest")
}
fn store_index_payload(
&self,
fingerprint: NtfsVolumeIndexFingerprint,
source_label: &'static str,
usn_checkpoint: Option<ScanCacheUsnCheckpoint>,
caveats: &[ParseCaveat],
mft_index: &MftIndex,
) -> Result<()> {
let payload = NtfsVolumeIndexPayload::new(
fingerprint.clone(),
source_label,
caveats.to_vec(),
mft_index.clone(),
);
let raw = serde_json::to_vec(&payload)?;
let payload_file = self.payload_file_for(&fingerprint);
write_ntfs_volume_index_file(&payload_file, &raw, "payload")?;
let payload_ref =
NtfsVolumeIndexPayloadRef::new(self.payload_file_name_for(&fingerprint), &raw);
let manifest = NtfsVolumeIndexCacheManifest::new(fingerprint, source_label, usn_checkpoint)
.with_payload(payload_ref);
self.store(&manifest)
}
}
fn validate_ntfs_volume_index_replay_range(
checkpoint: Option<&ScanCacheUsnCheckpoint>,
journal_state: &ScanCacheUsnJournalState,
) -> Option<NtfsVolumeIndexPayloadMiss> {
let Some(checkpoint) = checkpoint else {
return Some(NtfsVolumeIndexPayloadMiss::UsnCheckpointMissing);
};
if checkpoint.journal_id != journal_state.journal_id {
return Some(NtfsVolumeIndexPayloadMiss::UsnJournalChanged);
}
if journal_state.first_usn > checkpoint.next_usn || journal_state.next_usn < checkpoint.next_usn
{
return Some(NtfsVolumeIndexPayloadMiss::UsnRangeUnavailable);
}
None
}
fn persistent_cache_miss_caveat(reason: impl AsRef<str>) -> ParseCaveat {
ParseCaveat::new(
MFT_PERSISTENT_CACHE_MISS_CAVEAT_CODE,
format!(
"persistent NTFS/MFT volume-index cache missed; reason={}",
reason.as_ref()
),
)
}
fn persistent_cache_write_skipped_caveat(reason: impl AsRef<str>) -> ParseCaveat {
ParseCaveat::new(
MFT_PERSISTENT_CACHE_WRITE_SKIPPED_CAVEAT_CODE,
format!(
"persistent NTFS/MFT volume-index payload write skipped; reason={}",
reason.as_ref()
),
)
}
fn is_transient_persistent_cache_caveat(code: &str) -> bool {
matches!(
code,
MFT_PERSISTENT_CACHE_MISS_CAVEAT_CODE | MFT_PERSISTENT_CACHE_WRITE_SKIPPED_CAVEAT_CODE
)
}
fn stable_ntfs_volume_index_checkpoint(
before: Option<&ScanCacheUsnJournalState>,
after: Option<&ScanCacheUsnJournalState>,
) -> Option<ScanCacheUsnCheckpoint> {
let before = before?;
let after = after?;
if before.journal_id != after.journal_id || before.next_usn != after.next_usn {
return None;
}
if after.first_usn > after.next_usn {
return None;
}
Some(ScanCacheUsnCheckpoint {
journal_id: after.journal_id,
next_usn: after.next_usn,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference,
parent_reference: NtfsFileReference,
usn: u64,
}
fn validate_ntfs_volume_index_replay(
mft_index: &MftIndex,
target_record_id: u64,
target_is_volume_root: bool,
checkpoint_next_usn: u64,
journal_next_usn: u64,
changes: &[NtfsVolumeIndexUsnChange],
) -> Option<NtfsVolumeIndexPayloadMiss> {
if journal_next_usn <= checkpoint_next_usn {
return None;
}
let replayed_changes = changes
.iter()
.filter(|change| change.usn >= checkpoint_next_usn && change.usn < journal_next_usn)
.collect::<Vec<_>>();
if target_is_volume_root {
return if replayed_changes.is_empty() {
Some(NtfsVolumeIndexPayloadMiss::UsnRangeUnavailable)
} else {
Some(NtfsVolumeIndexPayloadMiss::UsnTargetChanged)
};
}
for change in replayed_changes {
match ntfs_usn_change_touches_target(mft_index, target_record_id, change) {
Ok(true) => return Some(NtfsVolumeIndexPayloadMiss::UsnTargetChanged),
Ok(false) => {}
Err(()) => return Some(NtfsVolumeIndexPayloadMiss::UsnAncestryUnavailable),
}
}
None
}
fn ntfs_usn_change_touches_target(
mft_index: &MftIndex,
target_record_id: u64,
change: &NtfsVolumeIndexUsnChange,
) -> std::result::Result<bool, ()> {
if change.file_reference.record_id == target_record_id
|| change.parent_reference.record_id == target_record_id
{
return Ok(true);
}
if ntfs_parent_chain_touches_target(mft_index, target_record_id, change.parent_reference)? {
return Ok(true);
}
if let Some(entry) = mft_index.get(change.file_reference.record_id) {
if entry.reference.sequence_number != change.file_reference.sequence_number {
return Err(());
}
for candidate in &entry.path_candidates {
if ntfs_parent_chain_touches_target(
mft_index,
target_record_id,
candidate.parent_reference,
)? {
return Ok(true);
}
}
}
Ok(false)
}
fn ntfs_parent_chain_touches_target(
mft_index: &MftIndex,
target_record_id: u64,
mut parent_reference: NtfsFileReference,
) -> std::result::Result<bool, ()> {
let mut visited = BTreeSet::new();
loop {
if parent_reference.record_id == target_record_id {
return Ok(true);
}
if !visited.insert(parent_reference.record_id) {
return Err(());
}
let Some(parent_entry) = mft_index.get(parent_reference.record_id) else {
return Err(());
};
if parent_entry.reference.sequence_number != parent_reference.sequence_number {
return Err(());
}
let next_parent = parent_entry.parent_reference;
if next_parent.record_id == parent_reference.record_id {
return Ok(false);
}
parent_reference = next_parent;
}
}
fn write_ntfs_volume_index_file(cache_file: &Path, raw: &[u8], label: &str) -> Result<()> {
let parent = cache_file.parent().ok_or_else(|| {
RebeccaError::ScanCacheUnavailable(format!(
"NTFS volume-index {label} path has no parent: {}",
cache_file.display()
))
})?;
std::fs::create_dir_all(parent).map_err(|err| {
RebeccaError::ScanCacheUnavailable(format!(
"NTFS volume-index {label} directory unavailable at {}: {}",
parent.display(),
err
))
})?;
let temp_file = temp_manifest_file(cache_file);
let write_result = (|| -> std::io::Result<()> {
let mut file = std::fs::File::create(&temp_file)?;
file.write_all(raw)?;
replace_manifest_file(&temp_file, cache_file)
})();
if let Err(err) = write_result {
let _ = std::fs::remove_file(&temp_file);
return Err(RebeccaError::ScanCacheUnavailable(format!(
"NTFS volume-index {label} write failed at {}: {}",
cache_file.display(),
err
)));
}
Ok(())
}
fn prune_manifest_file(cache_file: &Path) {
if let Err(err) = std::fs::remove_file(cache_file)
&& err.kind() != std::io::ErrorKind::NotFound
{
tracing::debug!(
path = %cache_file.display(),
error = %err,
"NTFS volume-index manifest prune skipped"
);
}
}
fn prune_manifest_pair(manifest_file: &Path, payload_file: &Path) {
prune_manifest_file(manifest_file);
prune_manifest_file(payload_file);
}
fn replace_manifest_file(temp_file: &Path, cache_file: &Path) -> std::io::Result<()> {
match std::fs::rename(temp_file, cache_file) {
Ok(()) => Ok(()),
Err(_) if cache_file.exists() => {
std::fs::remove_file(cache_file)?;
std::fs::rename(temp_file, cache_file)
}
Err(err) => Err(err),
}
}
fn temp_manifest_file(cache_file: &Path) -> PathBuf {
let file_name = cache_file
.file_name()
.map(|name| name.to_string_lossy())
.unwrap_or_else(|| "ntfs-volume-index.json".into());
let unique = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_nanos())
.unwrap_or_default();
cache_file.with_file_name(format!("{file_name}.tmp-{}-{unique}", std::process::id()))
}
fn unix_now() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_secs())
.unwrap_or_default()
}
#[derive(Debug, Default)]
pub(super) struct WindowsNtfsMftIndexCache {
volumes: Mutex<BTreeMap<NtfsVolumeIndexCacheKey, CachedNtfsVolumeIndexSlot>>,
manifest_store: Option<NtfsVolumeIndexManifestStore>,
volume_changed: Condvar,
}
#[derive(Debug, Default)]
struct PersistentIndexLoad {
index: Option<CachedNtfsVolumeIndex>,
caveats: Vec<ParseCaveat>,
backend_evidence: ScanBackendEvidence,
}
impl PersistentIndexLoad {
fn hit(index: CachedNtfsVolumeIndex) -> Self {
let mut backend_evidence = ScanBackendEvidence::default();
backend_evidence.record_cache_event("ntfs-volume-index", "hit", None);
Self {
index: Some(index),
caveats: Vec::new(),
backend_evidence,
}
}
fn miss(reason: impl AsRef<str>) -> Self {
let reason = reason.as_ref();
let mut backend_evidence = ScanBackendEvidence::default();
backend_evidence.record_cache_event("ntfs-volume-index", "miss", Some(reason.to_string()));
Self {
index: None,
caveats: vec![persistent_cache_miss_caveat(reason)],
backend_evidence,
}
}
fn disabled() -> Self {
Self::default()
}
}
impl WindowsNtfsMftIndexCache {
pub(super) fn with_manifest_store(manifest_store: NtfsVolumeIndexManifestStore) -> Self {
Self {
volumes: Mutex::default(),
manifest_store: Some(manifest_store),
volume_changed: Condvar::new(),
}
}
fn load_or_build<'a>(
&self,
capabilities: &NtfsVolumeCapabilities,
target_record_id: u64,
target_is_volume_root: bool,
cancellation: &ScanCancellationToken,
observer: Option<&'a mut NtfsMftBuildObserver<'a>>,
) -> Result<Arc<CachedNtfsVolumeIndex>> {
let cache_key = capabilities.cache_key();
let mut volumes = self.lock_volumes()?;
loop {
check_not_cancelled(cancellation)?;
match volumes.get(&cache_key) {
Some(CachedNtfsVolumeIndexSlot::Ready(index))
if index.is_reusable_for(capabilities) =>
{
return Ok(Arc::clone(index));
}
Some(CachedNtfsVolumeIndexSlot::Ready(_)) => {
volumes.remove(&cache_key);
}
Some(CachedNtfsVolumeIndexSlot::Unavailable(reason)) => {
return Err(RebeccaError::PlatformUnavailable(reason.clone()));
}
Some(CachedNtfsVolumeIndexSlot::Building) => {
volumes = self.wait_for_volume_update(volumes)?;
}
None => {
volumes.insert(cache_key.clone(), CachedNtfsVolumeIndexSlot::Building);
break;
}
}
}
drop(volumes);
let persistent_load = self.load_persistent_index(
capabilities,
target_record_id,
target_is_volume_root,
cancellation,
)?;
let mut persistent_load_caveats = persistent_load.caveats;
let persistent_load_evidence = persistent_load.backend_evidence;
let build_result = match persistent_load.index {
Some(index) => Ok((index, Vec::new())),
None => CachedNtfsVolumeIndex::build(
capabilities,
cancellation,
self.manifest_store.is_some(),
observer,
)
.map(|index| (index, std::mem::take(&mut persistent_load_caveats))),
};
let mut volumes = self.lock_volumes()?;
let result = match build_result {
Ok((mut index, mut transient_caveats)) => {
index.backend_evidence.merge(persistent_load_evidence);
if let Some((caveat, reason)) =
index.store_persistent_payload(self.manifest_store.as_ref())
{
index.caveats.push(caveat);
index.backend_evidence.record_cache_event(
"ntfs-volume-index",
"write-skipped",
Some(reason.to_string()),
);
}
index.caveats.append(&mut transient_caveats);
let index = Arc::new(index);
volumes.insert(
cache_key,
CachedNtfsVolumeIndexSlot::Ready(Arc::clone(&index)),
);
Ok(index)
}
Err(err) => {
if let Some(reason) = cacheable_index_failure(&err) {
volumes.insert(cache_key, CachedNtfsVolumeIndexSlot::Unavailable(reason));
} else {
volumes.remove(&cache_key);
}
Err(err)
}
};
self.volume_changed.notify_all();
result
}
fn load_persistent_index(
&self,
capabilities: &NtfsVolumeCapabilities,
target_record_id: u64,
target_is_volume_root: bool,
cancellation: &ScanCancellationToken,
) -> Result<PersistentIndexLoad> {
let Some(manifest_store) = self.manifest_store.as_ref() else {
return Ok(PersistentIndexLoad::disabled());
};
check_not_cancelled(cancellation)?;
let volume = match LiveNtfsVolume::open(capabilities) {
Ok(volume) => volume,
Err(err) => {
tracing::debug!(
error = %err,
"NTFS persistent volume-index load skipped before open"
);
return Ok(PersistentIndexLoad::miss("open-failed"));
}
};
let volume_data = match volume.ntfs_volume_data() {
Ok(volume_data) => volume_data,
Err(err) => {
tracing::debug!(
error = %err,
"NTFS persistent volume-index load skipped before volume fingerprint"
);
return Ok(PersistentIndexLoad::miss("volume-data-unavailable"));
}
};
let geometry = match NtfsRecordGeometry::from_volume_data(&volume.device_path, &volume_data)
{
Ok(geometry) => geometry,
Err(err) => {
tracing::debug!(
error = %err,
"NTFS persistent volume-index load skipped before record geometry"
);
return Ok(PersistentIndexLoad::miss("record-geometry-unavailable"));
}
};
let fingerprint =
NtfsVolumeIndexFingerprint::from_volume_data(capabilities, &volume_data, geometry);
let journal_state = match volume.usn_journal_state() {
Ok(journal_state) => journal_state,
Err(err) => {
tracing::debug!(
error = %err,
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index load skipped before USN freshness validation"
);
return Ok(PersistentIndexLoad::miss("usn-journal-unavailable"));
}
};
match manifest_store.load_replayable_index_payload(&fingerprint, &journal_state) {
NtfsVolumeIndexPayloadLookup::Hit { manifest, payload } => {
let manifest = *manifest;
let payload = *payload;
let Some(usn_checkpoint) = manifest.usn_checkpoint.clone() else {
tracing::debug!(
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index payload miss without USN checkpoint"
);
return Ok(PersistentIndexLoad::miss("usn-checkpoint-missing"));
};
let mut journal_state = journal_state;
let mut replay_attempts = 0_usize;
loop {
if journal_state.next_usn > usn_checkpoint.next_usn {
let changes = match volume.read_usn_changes(
&usn_checkpoint,
&journal_state,
cancellation,
) {
Ok(changes) => changes,
Err(err) => {
tracing::debug!(
error = %err,
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index payload miss because USN replay failed"
);
return Ok(PersistentIndexLoad::miss("usn-replay-read-failed"));
}
};
if let Some(reason) = validate_ntfs_volume_index_replay(
&payload.mft_index,
target_record_id,
target_is_volume_root,
usn_checkpoint.next_usn,
journal_state.next_usn,
&changes,
) {
tracing::debug!(
?reason,
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index payload miss after USN replay"
);
return Ok(PersistentIndexLoad::miss(reason.label()));
}
}
let replay_end_state = match volume.usn_journal_state() {
Ok(journal_state) => journal_state,
Err(err) => {
tracing::debug!(
error = %err,
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index payload miss after replay stability check failed"
);
return Ok(PersistentIndexLoad::miss(
"usn-replay-stability-check-failed",
));
}
};
if replay_end_state == journal_state {
break;
}
if let Some(reason) = validate_ntfs_volume_index_replay_range(
Some(&usn_checkpoint),
&replay_end_state,
) {
tracing::debug!(
?reason,
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index payload miss after USN replay state changed"
);
return Ok(PersistentIndexLoad::miss(reason.label()));
}
replay_attempts += 1;
if replay_attempts >= MAX_NTFS_VOLUME_INDEX_USN_REPLAY_ATTEMPTS {
tracing::debug!(
reason = ?NtfsVolumeIndexPayloadMiss::UsnReplayUnstable,
generation = fingerprint.persistent_generation(),
attempts = replay_attempts,
"NTFS persistent volume-index payload miss because USN replay did not stabilize"
);
return Ok(PersistentIndexLoad::miss(
NtfsVolumeIndexPayloadMiss::UsnReplayUnstable.label(),
));
}
journal_state = replay_end_state;
}
tracing::debug!(
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index payload hit"
);
Ok(PersistentIndexLoad::hit(CachedNtfsVolumeIndex {
fingerprint,
mft_index: payload.mft_index,
source_label: PERSISTENT_MFT_SOURCE_LABEL,
caveats: payload.caveats,
backend_evidence: ScanBackendEvidence::default(),
usn_checkpoint: Some(usn_checkpoint),
}))
}
NtfsVolumeIndexPayloadLookup::Miss(reason) => {
tracing::debug!(
?reason,
generation = fingerprint.persistent_generation(),
"NTFS persistent volume-index payload miss"
);
Ok(PersistentIndexLoad::miss(reason.label()))
}
}
}
fn lock_volumes(
&self,
) -> Result<
std::sync::MutexGuard<'_, BTreeMap<NtfsVolumeIndexCacheKey, CachedNtfsVolumeIndexSlot>>,
> {
self.volumes.lock().map_err(|_| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} volume index cache is unavailable"
))
})
}
fn wait_for_volume_update<'a>(
&self,
volumes: std::sync::MutexGuard<
'a,
BTreeMap<NtfsVolumeIndexCacheKey, CachedNtfsVolumeIndexSlot>,
>,
) -> Result<
std::sync::MutexGuard<'a, BTreeMap<NtfsVolumeIndexCacheKey, CachedNtfsVolumeIndexSlot>>,
> {
let (volumes, _) = self
.volume_changed
.wait_timeout(volumes, Duration::from_millis(250))
.map_err(|_| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} volume index cache is unavailable"
))
})?;
Ok(volumes)
}
}
#[derive(Debug)]
enum CachedNtfsVolumeIndexSlot {
Ready(Arc<CachedNtfsVolumeIndex>),
Unavailable(String),
Building,
}
fn cacheable_index_failure(err: &RebeccaError) -> Option<String> {
match err {
RebeccaError::PlatformUnavailable(reason) => Some(reason.clone()),
_ => None,
}
}
#[derive(Debug)]
struct CachedNtfsVolumeIndex {
fingerprint: NtfsVolumeIndexFingerprint,
mft_index: MftIndex,
source_label: &'static str,
caveats: Vec<ParseCaveat>,
backend_evidence: ScanBackendEvidence,
usn_checkpoint: Option<ScanCacheUsnCheckpoint>,
}
impl CachedNtfsVolumeIndex {
fn is_reusable_for(&self, capabilities: &NtfsVolumeCapabilities) -> bool {
self.fingerprint.is_reusable_for(capabilities)
}
fn store_persistent_payload(
&self,
manifest_store: Option<&NtfsVolumeIndexManifestStore>,
) -> Option<(ParseCaveat, &'static str)> {
let manifest_store = manifest_store?;
let Some(usn_checkpoint) = self.usn_checkpoint.clone() else {
tracing::debug!(
generation = self.fingerprint.persistent_generation(),
"NTFS volume-index payload write skipped without stable USN checkpoint"
);
let reason = "stable-usn-checkpoint-unavailable";
return Some((persistent_cache_write_skipped_caveat(reason), reason));
};
if matches!(
manifest_store.load_index_payload(&self.fingerprint),
NtfsVolumeIndexPayloadLookup::Hit { manifest, .. }
if manifest.usn_checkpoint.as_ref() == Some(&usn_checkpoint)
) {
return None;
}
let persistent_caveats = self
.caveats
.iter()
.filter(|caveat| !is_transient_persistent_cache_caveat(&caveat.code))
.cloned()
.collect::<Vec<_>>();
if let Err(err) = manifest_store.store_index_payload(
self.fingerprint.clone(),
self.source_label,
Some(usn_checkpoint),
&persistent_caveats,
&self.mft_index,
) {
tracing::debug!(
error = %err,
generation = self.fingerprint.persistent_generation(),
"NTFS volume-index payload write skipped"
);
let reason = "write-failed";
return Some((persistent_cache_write_skipped_caveat(reason), reason));
}
None
}
fn build<'a>(
capabilities: &NtfsVolumeCapabilities,
cancellation: &ScanCancellationToken,
capture_usn_checkpoint: bool,
observer: Option<&'a mut NtfsMftBuildObserver<'a>>,
) -> Result<Self> {
let monitor = ntfs_mft_build_monitor(observer);
check_mft_build_progress(cancellation, &monitor)?;
let volume =
monitor.measure_checked(NtfsMftBuildStage::OpenVolume, cancellation, || {
LiveNtfsVolume::open(capabilities)
})?;
let volume_data =
monitor.measure_checked(NtfsMftBuildStage::ReadVolumeData, cancellation, || {
volume.ntfs_volume_data()
})?;
let geometry = NtfsRecordGeometry::from_volume_data(&volume.device_path, &volume_data)?;
let fingerprint =
NtfsVolumeIndexFingerprint::from_volume_data(capabilities, &volume_data, geometry);
let starting_usn_state = if capture_usn_checkpoint {
optional_ntfs_usn_journal_state(&volume, "before full-index build")
} else {
None
};
let records = volume.read_mft_records(&volume_data, cancellation, &monitor)?;
let source_label = records.source_label;
let mut stream_source = LiveNtfsIndexStreamSource {
volume: &volume,
cancellation,
monitor: &monitor,
};
let (mft_index, mut caveats) = build_mft_index_from_records(
records,
geometry,
&mut stream_source,
cancellation,
&monitor,
)?;
if let Some(caveat) = monitor.timing_caveat() {
caveats.push(caveat);
}
let ending_usn_state = if capture_usn_checkpoint {
optional_ntfs_usn_journal_state(&volume, "after full-index build")
} else {
None
};
let usn_checkpoint = stable_ntfs_volume_index_checkpoint(
starting_usn_state.as_ref(),
ending_usn_state.as_ref(),
);
Ok(Self {
fingerprint,
mft_index,
source_label,
caveats,
backend_evidence: monitor.evidence(),
usn_checkpoint,
})
}
}
fn optional_ntfs_usn_journal_state(
volume: &LiveNtfsVolume,
stage: &'static str,
) -> Option<ScanCacheUsnJournalState> {
match volume.usn_journal_state() {
Ok(state) => Some(state),
Err(err) => {
tracing::debug!(
error = %err,
stage,
"NTFS volume-index USN checkpoint capture skipped"
);
None
}
}
}
#[derive(Debug, Clone, Copy)]
pub(super) struct WindowsNtfsMftScanBackend<'a> {
cache: &'a WindowsNtfsMftIndexCache,
}
impl<'a> WindowsNtfsMftScanBackend<'a> {
pub(super) const fn new(cache: &'a WindowsNtfsMftIndexCache) -> Self {
Self { cache }
}
}
impl ScanBackend for WindowsNtfsMftScanBackend<'_> {
fn kind(&self) -> ScanBackendKind {
ScanBackendKind::WindowsNtfsMftExperimental
}
fn measure_path_with_progress<F>(
&self,
request: ScanRequest<'_>,
_progress: F,
) -> Result<MeasuredScan>
where
F: for<'a> FnMut(ScanProgressEvent<'a>),
{
check_not_cancelled(request.cancellation)?;
let metadata = root_metadata(request.path)?;
if is_reparse_like(&metadata) {
return Err(RebeccaError::SafetyBlocked(
"symlink or reparse point traversal is disabled".to_string(),
));
}
let capabilities = NtfsVolumeCapabilities::resolve(request.path)?;
let target_identity = FileIdentity::from_path(request.path)?;
if target_identity.volume_serial != capabilities.volume_serial {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} target volume identity changed while resolving {}",
request.path.display()
)));
}
let target_record_id = target_identity.file_reference.record_id;
let (summary, source_label, shared_caveats, backend_evidence) =
match build_targeted_mft_summary(
&capabilities,
target_identity.file_reference,
request.cancellation,
) {
Ok((summary, caveats, evidence)) => {
(summary, TARGETED_MFT_SOURCE_LABEL, caveats, evidence)
}
Err(err)
if live_ntfs_mft_full_index_fallback_enabled()
&& mft_record_source_error_can_fallback(&err) =>
{
let index = self.cache.load_or_build(
&capabilities,
target_record_id,
is_volume_root_path(request.path, &capabilities.root_path),
request.cancellation,
None,
)?;
let Some(_) = index.mft_index.get(target_record_id) else {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not map {} to MFT record {}",
request.path.display(),
target_record_id
)));
};
let mut summary = index.mft_index.aggregate_subtree(target_record_id);
summary.caveats.push(ParseCaveat::new(
"mft-targeted-full-index-fallback",
format!(
"targeted NTFS/MFT traversal was unavailable ({err}); full-volume MFT index fallback was enabled by {LIVE_NTFS_MFT_FULL_INDEX_FALLBACK_ENV}"
),
));
(
summary,
index.source_label,
index.caveats.clone(),
index.backend_evidence.clone(),
)
}
Err(err) => return Err(err),
};
let report = ScanReport {
bytes_scanned: summary.bytes,
files_scanned: summary.files,
directories_scanned: summary.directories,
};
let measured = MeasuredScan::exact(report, self.kind())
.with_backend_source(mft_backend_source_label(source_label))
.with_backend_evidence(backend_evidence);
Ok(with_bounded_mft_caveats(
measured,
shared_caveats.into_iter().chain(summary.caveats),
))
}
}
pub(super) fn inspect_disk_map_with_progress<'a, F>(
cache: &WindowsNtfsMftIndexCache,
path: &'a Path,
options: DiskMapBackendOptions,
cancellation: &ScanCancellationToken,
progress: &'a mut F,
) -> Result<DiskMapBackendReport>
where
F: for<'event> FnMut(InspectProgressEvent<'event>) -> InspectProgressResult + 'a,
{
let mut observer = |event| emit_mft_build_progress(path, event, progress);
inspect_disk_map_inner(cache, path, options, cancellation, Some(&mut observer))
}
fn inspect_disk_map_inner<'a>(
cache: &WindowsNtfsMftIndexCache,
path: &'a Path,
options: DiskMapBackendOptions,
cancellation: &ScanCancellationToken,
observer: Option<&'a mut NtfsMftBuildObserver<'a>>,
) -> Result<DiskMapBackendReport> {
check_not_cancelled(cancellation)?;
let metadata = root_metadata(path)?;
if is_reparse_like(&metadata) {
return Err(RebeccaError::SafetyBlocked(
"symlink or reparse point traversal is disabled".to_string(),
));
}
let capabilities = NtfsVolumeCapabilities::resolve(path)?;
let target_identity = FileIdentity::from_path(path)?;
if target_identity.volume_serial != capabilities.volume_serial {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} target volume identity changed while resolving {}",
path.display()
)));
}
let target_record_id = target_identity.file_reference.record_id;
if !is_volume_root_path(path, &capabilities.root_path)
&& !live_ntfs_mft_full_index_fallback_enabled()
{
return build_targeted_mft_disk_map(
&capabilities,
target_identity.file_reference,
path,
options,
cancellation,
observer,
);
}
let index = cache.load_or_build(
&capabilities,
target_record_id,
is_volume_root_path(path, &capabilities.root_path),
cancellation,
observer,
)?;
let target_entry = index
.mft_index
.get(target_record_id)
.cloned()
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not map {} to MFT record {}",
path.display(),
target_record_id
))
})?;
let mut top_entries = DiskMapTopEntries::new(
options.top_limit,
options.top_sort,
options.entry_filter.clone(),
);
let mut groups = options.group_collector();
let mut caveats = Vec::new();
let mut visited_directories = BTreeSet::new();
let max_depth = options.max_depth.unwrap_or(usize::MAX);
let backend_source = mft_backend_source_label(index.source_label);
let entry_provenance = EstimateProvenance::from_backend_confidence_and_source(
ScanBackendKind::WindowsNtfsMftExperimental,
ScanEstimateConfidence::Exact,
Some(backend_source.clone()),
);
let aggregate = if target_entry.is_directory {
caveats.extend(target_entry.caveats.clone());
extend_mft_directory_edge_caveats(&index.mft_index, target_record_id, &mut caveats);
visited_directories.insert(target_record_id);
let mut aggregate = PhysicalMetricsAccumulator::default();
for edge in index
.mft_index
.child_edges(target_record_id)
.cloned()
.collect::<Vec<_>>()
{
check_not_cancelled(cancellation)?;
let Some(child) = index.mft_index.get(edge.child.record_id).cloned() else {
continue;
};
let child_path = path.join(&edge.name);
let child_aggregate = collect_mft_disk_map_entry(
&index.mft_index,
path,
child_path,
child,
1,
max_depth,
&entry_provenance,
&mut visited_directories,
&mut caveats,
&mut top_entries,
&mut groups,
capabilities.volume_serial,
cancellation,
)?;
aggregate.absorb_child(child_aggregate);
}
aggregate
} else {
collect_mft_disk_map_entry(
&index.mft_index,
path,
path.to_path_buf(),
target_entry,
0,
max_depth,
&entry_provenance,
&mut visited_directories,
&mut caveats,
&mut top_entries,
&mut groups,
capabilities.volume_serial,
cancellation,
)?
};
let metrics = disk_map_metrics_from_physical(aggregate.into_metrics());
let measured = MeasuredScan::exact(
ScanReport {
bytes_scanned: metrics.logical_bytes,
files_scanned: metrics.files,
directories_scanned: metrics.directories,
},
ScanBackendKind::WindowsNtfsMftExperimental,
)
.with_backend_source(backend_source)
.with_backend_evidence(index.backend_evidence.clone());
let measured =
with_bounded_mft_caveats(measured, index.caveats.clone().into_iter().chain(caveats));
Ok(DiskMapBackendReport {
metrics,
top_entries: top_entries.into_sorted_entries(),
groups,
diagnostics: Vec::new(),
estimate_provenance: EstimateProvenance::from_measured_scan(&measured),
})
}
fn emit_mft_build_progress<F>(
root: &Path,
event: NtfsMftBuildMonitorEvent,
progress: &mut F,
) -> InspectProgressResult
where
F: for<'event> FnMut(InspectProgressEvent<'event>) -> InspectProgressResult,
{
let backend = ScanBackendKind::WindowsNtfsMftExperimental;
match event {
NtfsMftBuildMonitorEvent::StageStarted { stage } => {
progress(InspectProgressEvent::BackendStageStarted {
root,
backend,
stage,
})
}
NtfsMftBuildMonitorEvent::StageFinished { stage } => {
progress(InspectProgressEvent::BackendStageFinished {
root,
backend,
stage,
})
}
NtfsMftBuildMonitorEvent::Metric { metric, value } => {
progress(InspectProgressEvent::BackendMetric {
root,
backend,
metric,
value,
})
}
}
}
#[expect(
clippy::too_many_arguments,
reason = "recursive traversal carries bounded report state"
)]
fn collect_mft_disk_map_entry(
index: &MftIndex,
root: &Path,
path: PathBuf,
entry: MftIndexEntry,
depth: usize,
max_depth: usize,
estimate_provenance: &EstimateProvenance,
visited: &mut BTreeSet<u64>,
caveats: &mut Vec<ParseCaveat>,
top_entries: &mut DiskMapTopEntries,
groups: &mut DiskMapGroupCollector,
volume_serial_number: u64,
cancellation: &ScanCancellationToken,
) -> Result<PhysicalMetricsAccumulator> {
check_not_cancelled(cancellation)?;
caveats.extend(entry.caveats.clone());
if entry.is_reparse_point {
caveats.push(ParseCaveat::new(
"reparse-point-skipped",
format!("record {} is a reparse point", entry.reference.record_id),
));
return Ok(PhysicalMetricsAccumulator::default());
}
let mut aggregate = PhysicalMetricsAccumulator::default();
if entry.is_directory {
if !visited.insert(entry.reference.record_id) {
caveats.push(ParseCaveat::new(
"mft-index-cycle-skipped",
format!(
"directory record {} appeared more than once in a subtree",
entry.reference.record_id
),
));
return Ok(PhysicalMetricsAccumulator::default());
}
extend_mft_directory_edge_caveats(index, entry.reference.record_id, caveats);
aggregate.record_directory();
groups.record_directory();
} else {
aggregate.record_file_path(
entry.reference.record_id,
entry.logical_size,
entry.allocated_size,
);
}
if entry.is_directory {
for edge in index
.child_edges(entry.reference.record_id)
.cloned()
.collect::<Vec<_>>()
{
let Some(child) = index.get(edge.child.record_id).cloned() else {
caveats.push(ParseCaveat::new(
"missing-record",
format!(
"record {} is not present in the MFT index",
edge.child.record_id
),
));
continue;
};
let child_path = path.join(&edge.name);
let child_aggregate = collect_mft_disk_map_entry(
index,
root,
child_path,
child,
depth.saturating_add(1),
max_depth,
estimate_provenance,
visited,
caveats,
top_entries,
groups,
volume_serial_number,
cancellation,
)?;
aggregate.absorb_child(child_aggregate);
}
}
if !entry.is_directory {
groups.record_file(
&path,
depth,
entry.logical_size,
entry.allocated_size,
ntfs_filetime_to_system_time(Some(entry.modified_windows_filetime)),
DiskMapMetadataSemantics::with_file_identity(DiskMapFileIdentity::new(
volume_serial_number,
entry.reference.record_id,
)),
);
}
if depth <= max_depth {
let metrics = disk_map_metrics_from_physical(aggregate.metrics());
top_entries.push(DiskMapEntry {
path,
root: root.to_path_buf(),
kind: mft_disk_map_entry_kind(&entry),
depth,
logical_bytes: metrics.logical_bytes,
allocated_bytes: metrics.allocated_bytes,
unique_logical_bytes: metrics.unique_logical_bytes,
unique_allocated_bytes: metrics.unique_allocated_bytes,
files: metrics.files,
directories: metrics.directories,
estimate_source: EstimateSource::FreshScan,
estimate_provenance: estimate_provenance.clone(),
cleanup_advice: None,
});
}
Ok(aggregate)
}
fn extend_mft_directory_edge_caveats(
index: &MftIndex,
parent_record_id: u64,
caveats: &mut Vec<ParseCaveat>,
) {
for edge in index.directory_edges(parent_record_id) {
caveats.extend(edge.caveats.clone());
}
}
fn mft_disk_map_entry_kind(entry: &MftIndexEntry) -> DiskMapEntryKind {
if entry.is_directory {
DiskMapEntryKind::Directory
} else {
DiskMapEntryKind::File
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct NtfsVolumeCapabilities {
root_path: PathBuf,
device_path: String,
mft_data_path: String,
volume_serial: u64,
}
impl NtfsVolumeCapabilities {
fn resolve(path: &Path) -> Result<Self> {
let volume_paths = VolumePaths::from_path(path)?;
if drive_type(&volume_paths.root_path) != DRIVE_FIXED {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} only indexes local fixed NTFS volumes"
)));
}
let info = volume_information(&volume_paths.root_path)?;
if !info
.file_system_name
.eq_ignore_ascii_case(NTFS_FILE_SYSTEM_NAME)
{
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} requires NTFS; {} uses {}",
volume_paths.root_path.display(),
info.file_system_name
)));
}
Ok(Self {
root_path: volume_paths.root_path,
device_path: volume_paths.device_path,
mft_data_path: volume_paths.mft_data_path,
volume_serial: u64::from(info.volume_serial),
})
}
fn cache_key(&self) -> NtfsVolumeIndexCacheKey {
NtfsVolumeIndexCacheKey::new(self.device_path.clone(), self.volume_serial)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct VolumePaths {
root_path: PathBuf,
device_path: String,
mft_data_path: String,
}
impl VolumePaths {
fn from_path(path: &Path) -> Result<Self> {
if !path.is_absolute() {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} requires an absolute local path"
)));
}
let Some(Component::Prefix(prefix)) = path.components().next() else {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not resolve a drive root for {}",
path.display()
)));
};
let drive = match prefix.kind() {
Prefix::Disk(drive) | Prefix::VerbatimDisk(drive) => drive,
Prefix::UNC(..)
| Prefix::VerbatimUNC(..)
| Prefix::DeviceNS(_)
| Prefix::Verbatim(_) => {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} does not index UNC or device namespace paths"
)));
}
};
let drive = char::from(drive).to_ascii_uppercase();
Ok(Self {
root_path: PathBuf::from(format!("{drive}:\\")),
device_path: format!("\\\\.\\{drive}:"),
mft_data_path: format!("\\\\?\\{drive}:\\$MFT::$DATA"),
})
}
}
fn is_volume_root_path(path: &Path, expected_root: &Path) -> bool {
if path == expected_root {
return true;
}
let mut components = path.components();
let is_drive_prefix = matches!(
components.next(),
Some(Component::Prefix(prefix))
if matches!(prefix.kind(), Prefix::Disk(_) | Prefix::VerbatimDisk(_))
);
is_drive_prefix
&& matches!(components.next(), Some(Component::RootDir))
&& components.next().is_none()
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct VolumeInformation {
volume_serial: u32,
file_system_name: String,
}
fn volume_information(root_path: &Path) -> Result<VolumeInformation> {
let root = wide_null(root_path.as_os_str());
let mut volume_serial = 0_u32;
let mut file_system_name = [0_u16; 32];
unsafe {
GetVolumeInformationW(
PCWSTR(root.as_ptr()),
None,
Some(&mut volume_serial),
None,
None,
Some(&mut file_system_name),
)
}
.map_err(|err| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not inspect volume {}: {}",
root_path.display(),
err.message()
))
})?;
Ok(VolumeInformation {
volume_serial,
file_system_name: wide_buffer_to_string(&file_system_name),
})
}
fn drive_type(root_path: &Path) -> u32 {
let root = wide_null(root_path.as_os_str());
unsafe { GetDriveTypeW(PCWSTR(root.as_ptr())) }
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct FileIdentity {
volume_serial: u64,
file_reference: NtfsFileReference,
}
impl FileIdentity {
fn from_path(path: &Path) -> Result<Self> {
let file = std::fs::OpenOptions::new()
.access_mode(FILE_READ_ATTRIBUTES.0)
.share_mode(FILE_SHARE_READ.0 | FILE_SHARE_WRITE.0 | FILE_SHARE_DELETE.0)
.custom_flags(FILE_FLAG_BACKUP_SEMANTICS.0)
.open(path)
.map_err(|err| {
RebeccaError::ScanFailed(ScanFailure::from_io(
path,
ScanFailurePhase::RootMetadata,
&err,
))
})?;
let mut info = BY_HANDLE_FILE_INFORMATION::default();
unsafe { GetFileInformationByHandle(HANDLE(file.as_raw_handle()), &mut info) }.map_err(
|err| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not read file identity for {}: {}",
path.display(),
err.message()
))
},
)?;
Ok(Self {
volume_serial: u64::from(info.dwVolumeSerialNumber),
file_reference: file_reference_from_number(
(u64::from(info.nFileIndexHigh) << 32) | u64::from(info.nFileIndexLow),
),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ParsedNtfsRecords {
source_label: &'static str,
records: Vec<NtfsParsedRecord>,
caveats: Vec<ParseCaveat>,
}
trait MftRecordSource {
fn label(&self) -> &'static str;
fn read_records(
&self,
volume_data: &NTFS_VOLUME_DATA_BUFFER,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<ParsedNtfsRecords>;
}
fn read_mft_records_from_sources(
sources: &[&dyn MftRecordSource],
volume_data: &NTFS_VOLUME_DATA_BUFFER,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<ParsedNtfsRecords> {
let mut fallback_errors = Vec::new();
for source in sources {
check_mft_build_progress(cancellation, monitor)?;
match source.read_records(volume_data, cancellation, monitor) {
Ok(mut records) => {
monitor.add_metric_usize(NtfsMftBuildMetric::ParsedRecords, records.records.len());
records.source_label = source.label();
records.caveats.extend(
fallback_errors
.drain(..)
.map(|reason| ParseCaveat::new("mft-record-source-fallback", reason)),
);
return Ok(records);
}
Err(err) if monitor.is_timed_out() => return Err(err),
Err(err) if mft_record_source_error_can_fallback(&err) => {
fallback_errors.push(format!("{}: {err}", source.label()));
}
Err(err) => return Err(err),
}
}
Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} record sources are unavailable: {}",
fallback_errors.join("; ")
)))
}
fn mft_record_source_error_can_fallback(err: &RebeccaError) -> bool {
matches!(
err,
RebeccaError::PlatformUnavailable(_) | RebeccaError::ScanFailed(_)
)
}
fn mft_backend_source_label(source_label: &str) -> String {
format!("{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL}-{source_label}")
}
#[derive(Debug, Default)]
struct MftParseErrorCaveats {
total: usize,
samples: Vec<ParseCaveat>,
}
impl MftParseErrorCaveats {
fn record(&mut self, record_id: u64, error: impl fmt::Display) {
self.total = self.total.saturating_add(1);
if self.samples.len() < MAX_MFT_PARSE_ERROR_CAVEAT_SAMPLES {
self.samples.push(ParseCaveat::new(
"mft-record-parse-error",
format!("record {record_id} could not be parsed: {error}"),
));
}
}
fn append_to(self, caveats: &mut Vec<ParseCaveat>) {
if self.total == 0 {
return;
}
let sample_count = self.samples.len();
caveats.extend(self.samples);
let omitted = self.total.saturating_sub(sample_count);
if omitted > 0 {
caveats.push(ParseCaveat::new(
"mft-record-parse-error-summary",
format!(
"{omitted} additional MFT records could not be parsed; parse-error samples were capped at {sample_count}"
),
));
}
}
}
#[derive(Debug, Default)]
struct BoundedMftCaveatBucket {
total: usize,
samples: Vec<String>,
}
fn with_bounded_mft_caveats<I>(mut measured: MeasuredScan, caveats: I) -> MeasuredScan
where
I: IntoIterator<Item = ParseCaveat>,
{
let mut buckets: BTreeMap<String, BoundedMftCaveatBucket> = BTreeMap::new();
for caveat in caveats {
let bucket = buckets.entry(caveat.code).or_default();
bucket.total = bucket.total.saturating_add(1);
if bucket.samples.len() < MAX_MFT_ESTIMATE_CAVEAT_SAMPLES_PER_CODE {
bucket.samples.push(caveat.message);
}
}
for (code, bucket) in buckets {
let sample_count = bucket.samples.len();
let omitted = bucket.total.saturating_sub(sample_count);
for message in bucket.samples {
measured = measured.with_caveat(code.clone(), message);
}
if omitted > 0 {
measured = measured.with_caveat(
MFT_CAVEAT_SUMMARY_CODE,
format!(
"{omitted} additional '{code}' caveats were omitted from this estimate; samples are capped at {MAX_MFT_ESTIMATE_CAVEAT_SAMPLES_PER_CODE} per caveat code"
),
);
}
}
measured
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct NtfsRecordGeometry {
record_size: usize,
sector_size: usize,
bytes_per_cluster: u64,
max_record_count: u64,
}
impl NtfsRecordGeometry {
fn from_volume_data(device_path: &str, volume_data: &NTFS_VOLUME_DATA_BUFFER) -> Result<Self> {
let record_size = usize::try_from(volume_data.BytesPerFileRecordSegment).unwrap_or(0);
let sector_size = usize::try_from(volume_data.BytesPerSector).unwrap_or(0);
let bytes_per_cluster = u64::from(volume_data.BytesPerCluster);
let mft_valid_data_length = u64::try_from(volume_data.MftValidDataLength).unwrap_or(0);
if record_size == 0 || sector_size == 0 || bytes_per_cluster == 0 {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} received invalid NTFS record geometry from {device_path}"
)));
}
if !record_size.is_multiple_of(sector_size) {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} received unaligned NTFS record geometry from {device_path}"
)));
}
if mft_valid_data_length == 0 {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} received empty NTFS MFT metadata from {device_path}"
)));
}
let max_record_count = mft_valid_data_length.saturating_div(record_size as u64);
if max_record_count == 0 {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} NTFS MFT metadata is smaller than one file record on {device_path}"
)));
}
Ok(Self {
record_size,
sector_size,
bytes_per_cluster,
max_record_count,
})
}
fn stream_geometry(self) -> NtfsStreamGeometry {
NtfsStreamGeometry::new(self.bytes_per_cluster, self.sector_size)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct MftMirrorReadPlan {
base_record_id: u64,
volume_offset: u64,
byte_len: usize,
record_count: u64,
}
fn mft_mirror_read_plan(
device_path: &str,
volume_data: &NTFS_VOLUME_DATA_BUFFER,
geometry: NtfsRecordGeometry,
) -> Result<Option<MftMirrorReadPlan>> {
let mirror_start_lcn = u64::try_from(volume_data.Mft2StartLcn).unwrap_or(0);
if mirror_start_lcn == 0 {
return Ok(None);
}
let record_count = MFT_MIRROR_SYSTEM_RECORDS.min(geometry.max_record_count);
if record_count == 0 {
return Ok(None);
}
let volume_offset = mirror_start_lcn
.checked_mul(geometry.bytes_per_cluster)
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} $MFTMirr offset overflowed on {device_path}"
))
})?;
let byte_len = record_count
.checked_mul(geometry.record_size as u64)
.and_then(|bytes| usize::try_from(bytes).ok())
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} $MFTMirr read length overflowed on {device_path}"
))
})?;
Ok(Some(MftMirrorReadPlan {
base_record_id: 0,
volume_offset,
byte_len,
record_count,
}))
}
fn build_mft_index_from_records<S>(
records: ParsedNtfsRecords,
geometry: NtfsRecordGeometry,
source: &mut S,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<(MftIndex, Vec<ParseCaveat>)>
where
S: NtfsStreamSource,
{
check_mft_build_progress(cancellation, monitor)?;
let mut caveats = records.caveats;
let record_set = monitor.measure_allow_budget_overrun_after_success(
NtfsMftBuildStage::ResolveIndexAllocations,
cancellation,
|| {
Ok(NtfsRecordSet::resolve_with_stream_source(
records.records,
geometry.stream_geometry(),
source,
))
},
)?;
let budget_exhausted_after_resolution = monitor.is_timed_out();
if budget_exhausted_after_resolution {
caveats.push(index_allocation_budget_exhausted_caveat(monitor));
} else {
check_mft_build_progress(cancellation, monitor)?;
}
let index = if budget_exhausted_after_resolution {
monitor.measure_cancellation_only(NtfsMftBuildStage::BuildMftIndex, cancellation, || {
Ok(MftIndex::from_record_set(record_set))
})?
} else {
monitor.measure_checked(NtfsMftBuildStage::BuildMftIndex, cancellation, || {
Ok(MftIndex::from_record_set(record_set))
})?
};
if budget_exhausted_after_resolution {
check_not_cancelled(cancellation)?;
} else {
check_mft_build_progress(cancellation, monitor)?;
}
Ok((index, caveats))
}
fn index_allocation_budget_exhausted_caveat(monitor: &NtfsMftBuildMonitor) -> ParseCaveat {
let budget = monitor
.budget
.timeout
.map(|timeout| format!("{}s", timeout.as_secs()))
.unwrap_or_else(|| "disabled".to_string());
let timings = monitor
.build_summary()
.map(|summary| format!("; {summary}"))
.unwrap_or_default();
ParseCaveat::new(
MFT_INDEX_ALLOCATION_BUDGET_EXHAUSTED_CAVEAT_CODE,
format!(
"live NTFS/MFT index allocation expansion crossed the {budget} build budget after parsed MFT records were available; returning degraded full-index evidence from available records{timings}"
),
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TargetedMftTraversalLimits {
max_records: usize,
max_depth: usize,
}
impl Default for TargetedMftTraversalLimits {
fn default() -> Self {
Self {
max_records: TARGETED_MFT_MAX_RECORDS,
max_depth: TARGETED_MFT_MAX_DEPTH,
}
}
}
trait TargetedMftRecordResolver {
fn resolve_record(&mut self, reference: NtfsFileReference) -> Result<Option<NtfsParsedRecord>>;
}
fn build_targeted_mft_summary(
capabilities: &NtfsVolumeCapabilities,
target_reference: NtfsFileReference,
cancellation: &ScanCancellationToken,
) -> Result<(SubtreeSummary, Vec<ParseCaveat>, ScanBackendEvidence)> {
let monitor = NtfsMftBuildMonitor::from_environment();
check_mft_build_progress(cancellation, &monitor)?;
let volume = monitor.measure_checked(NtfsMftBuildStage::OpenVolume, cancellation, || {
LiveNtfsVolume::open(capabilities)
})?;
let volume_data =
monitor.measure_checked(NtfsMftBuildStage::ReadVolumeData, cancellation, || {
volume.ntfs_volume_data()
})?;
let geometry = NtfsRecordGeometry::from_volume_data(&volume.device_path, &volume_data)?;
let mut resolver = LiveNtfsTargetRecordResolver::new(&volume, geometry, cancellation, &monitor);
let mut stream_source = LiveNtfsIndexStreamSource {
volume: &volume,
cancellation,
monitor: &monitor,
};
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut stream_source,
geometry,
cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let summary = monitor.measure_checked(
NtfsMftBuildStage::TargetedTraverseSubtree,
cancellation,
|| traversal.aggregate_subtree(target_reference),
)?;
let mut caveats = Vec::new();
resolver.append_parse_caveats(&mut caveats);
if let Some(caveat) = monitor.timing_caveat() {
caveats.push(caveat);
}
Ok((summary, caveats, monitor.evidence()))
}
fn build_targeted_mft_disk_map<'a>(
capabilities: &NtfsVolumeCapabilities,
target_reference: NtfsFileReference,
root_path: &Path,
options: DiskMapBackendOptions,
cancellation: &ScanCancellationToken,
observer: Option<&'a mut NtfsMftBuildObserver<'a>>,
) -> Result<DiskMapBackendReport> {
let monitor = ntfs_mft_build_monitor(observer);
check_mft_build_progress(cancellation, &monitor)?;
let volume = monitor.measure_checked(NtfsMftBuildStage::OpenVolume, cancellation, || {
LiveNtfsVolume::open(capabilities)
})?;
let volume_data =
monitor.measure_checked(NtfsMftBuildStage::ReadVolumeData, cancellation, || {
volume.ntfs_volume_data()
})?;
let geometry = NtfsRecordGeometry::from_volume_data(&volume.device_path, &volume_data)?;
let mut resolver = LiveNtfsTargetRecordResolver::new(&volume, geometry, cancellation, &monitor);
let mut stream_source = LiveNtfsIndexStreamSource {
volume: &volume,
cancellation,
monitor: &monitor,
};
let entry_provenance = EstimateProvenance::from_backend_confidence_and_source(
ScanBackendKind::WindowsNtfsMftExperimental,
ScanEstimateConfidence::Exact,
Some(mft_backend_source_label(TARGETED_MFT_SOURCE_LABEL)),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut stream_source,
geometry,
cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let targeted_map = monitor.measure_checked(
NtfsMftBuildStage::TargetedTraverseSubtree,
cancellation,
|| {
traversal.collect_disk_map(
target_reference,
root_path,
&options,
capabilities.volume_serial,
&entry_provenance,
)
},
)?;
let mut caveats = targeted_map.caveats.clone();
resolver.append_parse_caveats(&mut caveats);
if let Some(caveat) = monitor.timing_caveat() {
caveats.push(caveat);
}
let measured = MeasuredScan::exact(
ScanReport {
bytes_scanned: targeted_map.metrics.logical_bytes,
files_scanned: targeted_map.metrics.files,
directories_scanned: targeted_map.metrics.directories,
},
ScanBackendKind::WindowsNtfsMftExperimental,
)
.with_backend_source(mft_backend_source_label(TARGETED_MFT_SOURCE_LABEL))
.with_backend_evidence(monitor.evidence());
let measured = with_bounded_mft_caveats(measured, caveats);
Ok(DiskMapBackendReport {
metrics: targeted_map.metrics,
top_entries: targeted_map.top_entries,
groups: targeted_map.groups,
diagnostics: Vec::new(),
estimate_provenance: EstimateProvenance::from_measured_scan(&measured),
})
}
struct LiveNtfsTargetRecordResolver<'a, 'observer> {
volume: &'a LiveNtfsVolume,
geometry: NtfsRecordGeometry,
cancellation: &'a ScanCancellationToken,
monitor: &'a NtfsMftBuildMonitor<'observer>,
records: BTreeMap<u64, NtfsParsedRecord>,
parse_errors: MftParseErrorCaveats,
}
impl<'a, 'observer> LiveNtfsTargetRecordResolver<'a, 'observer> {
fn new(
volume: &'a LiveNtfsVolume,
geometry: NtfsRecordGeometry,
cancellation: &'a ScanCancellationToken,
monitor: &'a NtfsMftBuildMonitor<'observer>,
) -> Self {
Self {
volume,
geometry,
cancellation,
monitor,
records: BTreeMap::new(),
parse_errors: MftParseErrorCaveats::default(),
}
}
fn append_parse_caveats(self, caveats: &mut Vec<ParseCaveat>) {
self.parse_errors.append_to(caveats);
}
fn read_targeted_file_record(
&self,
reference: NtfsFileReference,
) -> Result<Option<(u64, Vec<u8>)>> {
self.monitor
.add_metric(NtfsMftBuildMetric::TargetedRecordAttempts, 1);
match self
.volume
.read_file_record(file_reference_number(reference), self.geometry.record_size)
{
Ok(Some(record)) => {
self.monitor
.add_metric(NtfsMftBuildMetric::TargetedRecordSuccesses, 1);
Ok(Some(record))
}
Ok(None) => Ok(None),
Err(err) if windows_error_matches(&err, ERROR_HANDLE_EOF) => Ok(None),
Err(err) if windows_error_matches(&err, ERROR_INVALID_PARAMETER) => Ok(None),
Err(err) => Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {TARGETED_MFT_SOURCE_LABEL} could not read MFT record {} from {}: {}",
reference.record_id,
self.volume.device_path,
err.message()
))),
}
}
}
impl TargetedMftRecordResolver for LiveNtfsTargetRecordResolver<'_, '_> {
fn resolve_record(&mut self, reference: NtfsFileReference) -> Result<Option<NtfsParsedRecord>> {
check_mft_build_progress(self.cancellation, self.monitor)?;
if let Some(record) = self.records.get(&reference.record_id) {
return Ok(Some(record.clone()));
}
if reference.record_id >= self.geometry.max_record_count {
return Ok(None);
}
let read = self.monitor.measure_checked(
NtfsMftBuildStage::TargetedReadRecord,
self.cancellation,
|| self.read_targeted_file_record(reference),
)?;
let Some((record_id, raw_record)) = read else {
return Ok(None);
};
let parsed_record_id = low_file_reference_number(record_id);
if parsed_record_id != reference.record_id {
return Ok(None);
}
match NtfsParsedRecord::parse_fsctl_file_record(
parsed_record_id,
&raw_record,
self.geometry.sector_size,
) {
Ok(record) => {
self.monitor
.add_metric(NtfsMftBuildMetric::ParsedRecords, 1);
self.records.insert(parsed_record_id, record.clone());
Ok(Some(record))
}
Err(err) => {
let message = err.to_string();
let record_len = raw_record.len();
let signature = record_signature_hex(&raw_record);
self.parse_errors.record(parsed_record_id, err);
Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {TARGETED_MFT_SOURCE_LABEL} could not parse targeted MFT record {parsed_record_id} ({record_len} bytes, signature {signature}): {message}"
)))
}
}
}
}
fn record_signature_hex(record: &[u8]) -> String {
record
.iter()
.take(4)
.map(|byte| format!("{byte:02X}"))
.collect::<Vec<_>>()
.join("")
}
struct TargetedMftTraversal<'a, 'observer, R, S>
where
R: TargetedMftRecordResolver,
S: NtfsStreamSource,
{
resolver: &'a mut R,
stream_source: &'a mut S,
geometry: NtfsRecordGeometry,
cancellation: &'a ScanCancellationToken,
monitor: &'a NtfsMftBuildMonitor<'observer>,
limits: TargetedMftTraversalLimits,
}
impl<R, S> TargetedMftTraversal<'_, '_, R, S>
where
R: TargetedMftRecordResolver,
S: NtfsStreamSource,
{
fn aggregate_subtree(&mut self, root: NtfsFileReference) -> Result<SubtreeSummary> {
let mut summary = SubtreeSummary::default();
let mut stack = vec![TargetedTraversalNode {
reference: root,
depth: 0,
directory_entry: None,
}];
let mut visited = std::collections::BTreeSet::new();
let mut traversal_attempts = 0_usize;
while let Some(node) = stack.pop() {
check_mft_build_progress(self.cancellation, self.monitor)?;
if traversal_attempts >= self.limits.max_records {
return Err(targeted_mft_unavailable(format!(
"targeted traversal exceeded the {} record candidate budget",
self.limits.max_records
)));
}
traversal_attempts = traversal_attempts.saturating_add(1);
let Some(record) = self.resolver.resolve_record(node.reference)? else {
if node.directory_entry.is_none() {
return Err(targeted_mft_unavailable(format!(
"target root record {} could not be resolved",
node.reference.record_id
)));
}
summary.caveats.push(ParseCaveat::new(
"missing-record",
format!(
"record {} is not present in the targeted MFT traversal",
node.reference.record_id
),
));
continue;
};
if reference_sequence_mismatches(node.reference, record.reference) {
summary.caveats.push(ParseCaveat::new(
"directory-index-child-sequence-mismatch",
format!(
"targeted traversal expected record {} sequence {:?}, but current sequence is {:?}",
node.reference.record_id,
node.reference.sequence_number,
record.reference.sequence_number
),
));
continue;
}
if !visited.insert(record.reference.record_id) {
summary.caveats.push(ParseCaveat::new(
"mft-targeted-record-already-counted",
format!(
"record {} appeared more than once in targeted subtree traversal",
record.reference.record_id
),
));
continue;
}
let record = self.resolve_record(record)?;
summary.caveats.extend(record.caveats.clone());
if let Some(directory_entry) = &node.directory_entry {
self.push_directory_entry_parent_caveat(&record, directory_entry, &mut summary);
}
if record.non_dos_file_name_count() > 1 {
summary.caveats.push(ParseCaveat::new(
"hardlink-path-candidates",
format!(
"record {} has multiple non-DOS file names; targeted traversal counted the record once",
record.reference.record_id
),
));
}
if !record.in_use {
continue;
}
if record.is_reparse_point {
summary.caveats.push(ParseCaveat::new(
"reparse-point-skipped",
format!("record {} is a reparse point", record.reference.record_id),
));
continue;
}
if record.is_directory {
summary.directories = summary.directories.saturating_add(1);
if node.depth >= self.limits.max_depth && !record.directory_entries.is_empty() {
return Err(targeted_mft_unavailable(format!(
"targeted traversal reached depth {} below directory record {}",
node.depth, record.reference.record_id
)));
}
self.push_directory_children(&record, node.depth, &mut stack, &mut summary);
} else {
let files_before = summary.files;
summary.files = summary.files.saturating_add(1);
summary.bytes = summary.bytes.saturating_add(record.cleanup_logical_size());
summary.allocated_bytes = add_file_allocated_bytes(
summary.allocated_bytes,
files_before,
record.cleanup_allocated_size(),
);
}
}
Ok(summary)
}
fn collect_disk_map(
&mut self,
root: NtfsFileReference,
root_path: &Path,
options: &DiskMapBackendOptions,
volume_serial_number: u64,
entry_provenance: &EstimateProvenance,
) -> Result<TargetedDiskMap> {
let mut state = TargetedDiskMapState::new(options);
let root_node = TargetedDiskMapNode {
reference: root,
path: root_path.to_path_buf(),
depth: 0,
directory_entry: None,
};
let context = TargetedDiskMapContext {
root_path,
max_visible_depth: options.max_depth.unwrap_or(usize::MAX),
volume_serial_number,
entry_provenance,
};
let aggregate = self.collect_disk_map_record(root_node, false, &context, &mut state)?;
Ok(TargetedDiskMap {
metrics: disk_map_metrics_from_physical(aggregate.into_metrics()),
top_entries: state.top_entries.into_sorted_entries(),
groups: state.groups,
caveats: state.caveats,
})
}
fn collect_disk_map_record(
&mut self,
node: TargetedDiskMapNode,
include_root_directory: bool,
context: &TargetedDiskMapContext<'_>,
state: &mut TargetedDiskMapState,
) -> Result<PhysicalMetricsAccumulator> {
check_mft_build_progress(self.cancellation, self.monitor)?;
state.record_attempt(self.limits.max_records)?;
let Some(record) = self.resolver.resolve_record(node.reference)? else {
if node.directory_entry.is_none() {
return Err(targeted_mft_unavailable(format!(
"target root record {} could not be resolved",
node.reference.record_id
)));
}
state.caveats.push(ParseCaveat::new(
"missing-record",
format!(
"record {} is not present in the targeted MFT traversal",
node.reference.record_id
),
));
return Ok(PhysicalMetricsAccumulator::default());
};
if reference_sequence_mismatches(node.reference, record.reference) {
state.caveats.push(ParseCaveat::new(
"directory-index-child-sequence-mismatch",
format!(
"targeted traversal expected record {} sequence {:?}, but current sequence is {:?}",
node.reference.record_id,
node.reference.sequence_number,
record.reference.sequence_number
),
));
return Ok(PhysicalMetricsAccumulator::default());
}
let record = self.resolve_record(record)?;
state.caveats.extend(record.caveats.clone());
if let Some(directory_entry) = &node.directory_entry {
self.push_directory_entry_parent_caveat_to(
&record,
directory_entry,
&mut state.caveats,
);
}
if record.non_dos_file_name_count() > 1 {
state.caveats.push(ParseCaveat::new(
"hardlink-path-candidates",
format!(
"record {} has multiple non-DOS file names; targeted disk-map path metrics preserve visible names and unique metrics count the physical record once",
record.reference.record_id
),
));
}
if !record.in_use {
return Ok(PhysicalMetricsAccumulator::default());
}
if record.is_reparse_point {
state.caveats.push(ParseCaveat::new(
"reparse-point-skipped",
format!("record {} is a reparse point", record.reference.record_id),
));
return Ok(PhysicalMetricsAccumulator::default());
}
let mut aggregate = PhysicalMetricsAccumulator::default();
if record.is_directory {
if !state.visited_directories.insert(record.reference.record_id) {
state.caveats.push(ParseCaveat::new(
"mft-targeted-directory-already-visited",
format!(
"directory record {} appeared more than once in targeted disk-map traversal",
record.reference.record_id
),
));
return Ok(PhysicalMetricsAccumulator::default());
}
if node.depth >= self.limits.max_depth && !record.directory_entries.is_empty() {
return Err(targeted_mft_unavailable(format!(
"targeted traversal reached depth {} below directory record {}",
node.depth, record.reference.record_id
)));
}
if include_root_directory || node.directory_entry.is_some() {
aggregate.record_directory();
state.groups.record_directory();
}
self.collect_disk_map_children(
&record,
&node.path,
node.depth,
context,
state,
&mut aggregate,
)?;
} else {
aggregate.record_file_path(
record.reference.record_id,
record.cleanup_logical_size(),
record.cleanup_allocated_size(),
);
}
if !record.is_directory {
state.groups.record_file(
&node.path,
node.depth,
record.cleanup_logical_size(),
record.cleanup_allocated_size(),
ntfs_filetime_to_system_time(
record
.primary_file_name()
.map(|file_name| file_name.modified_windows_filetime),
),
DiskMapMetadataSemantics::with_file_identity(DiskMapFileIdentity::new(
context.volume_serial_number,
record.reference.record_id,
)),
);
}
let should_push_entry =
!record.is_directory || include_root_directory || node.directory_entry.is_some();
if should_push_entry && node.depth <= context.max_visible_depth {
let metrics = disk_map_metrics_from_physical(aggregate.metrics());
state.top_entries.push(DiskMapEntry {
path: node.path,
root: context.root_path.to_path_buf(),
kind: if record.is_directory {
DiskMapEntryKind::Directory
} else {
DiskMapEntryKind::File
},
depth: node.depth,
logical_bytes: metrics.logical_bytes,
allocated_bytes: metrics.allocated_bytes,
unique_logical_bytes: metrics.unique_logical_bytes,
unique_allocated_bytes: metrics.unique_allocated_bytes,
files: metrics.files,
directories: metrics.directories,
estimate_source: EstimateSource::FreshScan,
estimate_provenance: context.entry_provenance.clone(),
cleanup_advice: None,
});
}
Ok(aggregate)
}
fn collect_disk_map_children(
&mut self,
record: &NtfsParsedRecord,
parent_path: &Path,
parent_depth: usize,
context: &TargetedDiskMapContext<'_>,
state: &mut TargetedDiskMapState,
aggregate: &mut PhysicalMetricsAccumulator,
) -> Result<()> {
for entry in &record.directory_entries {
if is_dos_directory_entry(entry) {
continue;
}
if entry.parent.record_id != record.reference.record_id {
state.caveats.push(ParseCaveat::new(
"directory-index-parent-mismatch",
format!(
"$I30 entry '{}' declares parent {}, but was stored on directory {}",
entry.name, entry.parent.record_id, record.reference.record_id
),
));
continue;
}
if reference_sequence_mismatches(entry.parent, record.reference) {
state.caveats.push(ParseCaveat::new(
"parent-sequence-mismatch",
format!(
"$I30 entry '{}' references parent {} sequence {:?}, but current sequence is {:?}",
entry.name,
entry.parent.record_id,
entry.parent.sequence_number,
record.reference.sequence_number
),
));
continue;
}
let child = TargetedDiskMapNode {
reference: entry.child,
path: parent_path.join(&entry.name),
depth: parent_depth.saturating_add(1),
directory_entry: Some(entry.clone()),
};
let child_aggregate = self.collect_disk_map_record(child, true, context, state)?;
aggregate.absorb_child(child_aggregate);
}
Ok(())
}
fn push_directory_entry_parent_caveat(
&self,
record: &NtfsParsedRecord,
directory_entry: &NtfsDirectoryEntry,
summary: &mut SubtreeSummary,
) {
self.push_directory_entry_parent_caveat_to(record, directory_entry, &mut summary.caveats);
}
fn push_directory_entry_parent_caveat_to(
&self,
record: &NtfsParsedRecord,
directory_entry: &NtfsDirectoryEntry,
caveats: &mut Vec<ParseCaveat>,
) {
let parent_edge_exists = record.names.iter().any(|name| {
!matches!(name.namespace, rebecca_ntfs::FileNameNamespace::Dos)
&& name.parent == directory_entry.parent
&& name.name.eq_ignore_ascii_case(&directory_entry.name)
});
if !parent_edge_exists {
caveats.push(ParseCaveat::new(
"directory-index-parent-map-fallback",
format!(
"$I30 entry '{}' was used because it is not present in $FILE_NAME parent edges for directory {}",
directory_entry.name, directory_entry.parent.record_id
),
));
}
}
fn resolve_record(&mut self, record: NtfsParsedRecord) -> Result<NtfsParsedRecord> {
let resolver = &mut self.resolver;
self.monitor.measure_checked(
NtfsMftBuildStage::TargetedResolveRecord,
self.cancellation,
|| {
resolve_record_with_stream_source(
record,
self.geometry.stream_geometry(),
self.stream_source,
|reference| resolver.resolve_record(reference),
)
},
)
}
fn push_directory_children(
&mut self,
record: &NtfsParsedRecord,
depth: usize,
stack: &mut Vec<TargetedTraversalNode>,
summary: &mut SubtreeSummary,
) {
for entry in &record.directory_entries {
if is_dos_directory_entry(entry) {
continue;
}
if entry.parent.record_id != record.reference.record_id {
summary.caveats.push(ParseCaveat::new(
"directory-index-parent-mismatch",
format!(
"$I30 entry '{}' declares parent {}, but was stored on directory {}",
entry.name, entry.parent.record_id, record.reference.record_id
),
));
continue;
}
if reference_sequence_mismatches(entry.parent, record.reference) {
summary.caveats.push(ParseCaveat::new(
"parent-sequence-mismatch",
format!(
"$I30 entry '{}' references parent {} sequence {:?}, but current sequence is {:?}",
entry.name,
entry.parent.record_id,
entry.parent.sequence_number,
record.reference.sequence_number
),
));
continue;
}
stack.push(TargetedTraversalNode {
reference: entry.child,
depth: depth.saturating_add(1),
directory_entry: Some(entry.clone()),
});
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct TargetedTraversalNode {
reference: NtfsFileReference,
depth: usize,
directory_entry: Option<NtfsDirectoryEntry>,
}
#[derive(Debug)]
struct TargetedDiskMap {
metrics: DiskMapMetrics,
top_entries: Vec<DiskMapEntry>,
groups: DiskMapGroupCollector,
caveats: Vec<ParseCaveat>,
}
fn disk_map_metrics_from_physical(metrics: PhysicalMetrics) -> DiskMapMetrics {
DiskMapMetrics {
logical_bytes: metrics.logical_bytes,
allocated_bytes: metrics.allocated_bytes,
unique_logical_bytes: (metrics.files > 0).then_some(metrics.unique_logical_bytes),
unique_allocated_bytes: metrics.unique_allocated_bytes,
files: metrics.files,
directories: metrics.directories,
}
}
#[derive(Debug, Clone, Copy)]
struct TargetedDiskMapContext<'a> {
root_path: &'a Path,
max_visible_depth: usize,
volume_serial_number: u64,
entry_provenance: &'a EstimateProvenance,
}
#[derive(Debug)]
struct TargetedDiskMapState {
visited_directories: BTreeSet<u64>,
traversal_attempts: usize,
top_entries: DiskMapTopEntries,
groups: DiskMapGroupCollector,
caveats: Vec<ParseCaveat>,
}
impl TargetedDiskMapState {
fn new(options: &DiskMapBackendOptions) -> Self {
Self {
visited_directories: BTreeSet::new(),
traversal_attempts: 0,
top_entries: DiskMapTopEntries::new(
options.top_limit,
options.top_sort,
options.entry_filter.clone(),
),
groups: options.group_collector(),
caveats: Vec::new(),
}
}
fn record_attempt(&mut self, max_records: usize) -> Result<()> {
if self.traversal_attempts >= max_records {
return Err(targeted_mft_unavailable(format!(
"targeted traversal exceeded the {max_records} record candidate budget"
)));
}
self.traversal_attempts = self.traversal_attempts.saturating_add(1);
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct TargetedDiskMapNode {
reference: NtfsFileReference,
path: PathBuf,
depth: usize,
directory_entry: Option<NtfsDirectoryEntry>,
}
fn targeted_mft_unavailable(reason: String) -> RebeccaError {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {TARGETED_MFT_SOURCE_LABEL} {reason}"
))
}
fn reference_sequence_mismatches(expected: NtfsFileReference, actual: NtfsFileReference) -> bool {
if expected.record_id != actual.record_id {
return true;
}
matches!(
(expected.sequence_number, actual.sequence_number),
(Some(expected), Some(actual)) if expected != 0 && actual != 0 && expected != actual
)
}
fn is_dos_directory_entry(entry: &NtfsDirectoryEntry) -> bool {
matches!(entry.namespace, rebecca_ntfs::FileNameNamespace::Dos)
}
fn add_file_allocated_bytes(
current: Option<u64>,
files_before: u64,
file_allocated: Option<u64>,
) -> Option<u64> {
match (current, file_allocated) {
(None, Some(right)) if files_before == 0 => Some(right),
(Some(left), Some(right)) => Some(left.saturating_add(right)),
_ => None,
}
}
const WINDOWS_TICK_SECONDS: u64 = 10_000_000;
const WINDOWS_TO_UNIX_EPOCH_SECONDS: u64 = 11_644_473_600;
fn ntfs_filetime_to_system_time(filetime: Option<u64>) -> Option<SystemTime> {
let filetime = filetime?;
if filetime == 0 {
return None;
}
let seconds = filetime / WINDOWS_TICK_SECONDS;
let ticks = filetime % WINDOWS_TICK_SECONDS;
if seconds < WINDOWS_TO_UNIX_EPOCH_SECONDS {
return None;
}
Some(
UNIX_EPOCH
+ Duration::from_secs(seconds - WINDOWS_TO_UNIX_EPOCH_SECONDS)
+ Duration::from_nanos(ticks.saturating_mul(100)),
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct MftExtent {
starting_vcn: u64,
lcn: u64,
cluster_count: u64,
}
fn parse_retrieval_pointer_extents(buffer: &[u8]) -> Result<Vec<MftExtent>> {
let header_size = offset_of!(RETRIEVAL_POINTERS_BUFFER, Extents);
if buffer.len() < header_size {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} retrieval pointer buffer is truncated"
)));
}
let extent_count = unsafe {
ptr::read_unaligned(
buffer
.as_ptr()
.add(offset_of!(RETRIEVAL_POINTERS_BUFFER, ExtentCount))
.cast::<u32>(),
)
};
let extent_count = usize::try_from(extent_count).unwrap_or(usize::MAX);
if extent_count == 0 {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} $MFT retrieval pointer list is empty"
)));
}
let extents_size = extent_count
.checked_mul(size_of::<RETRIEVAL_POINTERS_BUFFER_0>())
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} retrieval pointer extent count overflowed"
))
})?;
let required_len = header_size.checked_add(extents_size).ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} retrieval pointer buffer length overflowed"
))
})?;
if buffer.len() < required_len {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} retrieval pointer buffer ended before all extents"
)));
}
let mut starting_vcn = unsafe {
ptr::read_unaligned(
buffer
.as_ptr()
.add(offset_of!(RETRIEVAL_POINTERS_BUFFER, StartingVcn))
.cast::<i64>(),
)
};
if starting_vcn < 0 {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} $MFT retrieval pointer list has a negative starting VCN"
)));
}
let mut extents = Vec::with_capacity(extent_count);
for index in 0..extent_count {
let offset = header_size + (index * size_of::<RETRIEVAL_POINTERS_BUFFER_0>());
let raw = unsafe {
ptr::read_unaligned(
buffer
.as_ptr()
.add(offset)
.cast::<RETRIEVAL_POINTERS_BUFFER_0>(),
)
};
if raw.NextVcn <= starting_vcn {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} $MFT retrieval pointer extent {index} is not ordered"
)));
}
if raw.Lcn < 0 {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} $MFT retrieval pointer extent {index} is sparse"
)));
}
extents.push(MftExtent {
starting_vcn: starting_vcn as u64,
lcn: raw.Lcn as u64,
cluster_count: (raw.NextVcn - starting_vcn) as u64,
});
starting_vcn = raw.NextVcn;
}
Ok(extents)
}
fn next_mft_chunk_len(
bytes_remaining_in_extent: u64,
records_remaining: u64,
record_size: usize,
) -> usize {
if record_size == 0 || records_remaining == 0 {
return 0;
}
let chunk_limit = SEQUENTIAL_MFT_CHUNK_BYTES.max(record_size) as u64;
let record_bytes_remaining = records_remaining.saturating_mul(record_size as u64);
let bytes_to_read = bytes_remaining_in_extent
.min(record_bytes_remaining)
.min(chunk_limit);
usize::try_from(bytes_to_read - (bytes_to_read % record_size as u64)).unwrap_or(0)
}
fn sequential_mft_parse_window_chunks() -> usize {
bounded_parallelism_budget().min(SEQUENTIAL_MFT_PARSE_WINDOW_CHUNKS)
}
fn run_scoped_mft_parse<R, F>(work: F) -> R
where
F: FnOnce() -> R + Send,
R: Send,
{
run_scoped_parallel_work(&MFT_PARSE_THREAD_POOL, "ntfs-mft-parse", work)
}
#[derive(Debug)]
struct SequentialMftChunk {
base_record_id: u64,
bytes: Vec<u8>,
mirror: Option<SequentialMftMirrorChunk>,
}
#[derive(Debug, Clone)]
struct SequentialMftMirrorChunk {
base_record_id: u64,
bytes: Vec<u8>,
}
fn parse_sequential_mft_chunks(
reader: &MftRecordReader,
cancellation: &ScanCancellationToken,
chunks: &[SequentialMftChunk],
) -> Result<Vec<MftRecordBatch>> {
run_scoped_mft_parse(|| {
chunks
.par_iter()
.map(|chunk| {
check_not_cancelled(cancellation)?;
Ok(match &chunk.mirror {
Some(mirror) => reader.parse_records_from_with_mirror(
chunk.base_record_id,
&chunk.bytes,
mirror.base_record_id,
&mirror.bytes,
),
None => reader.parse_records_from(chunk.base_record_id, &chunk.bytes),
})
})
.collect::<Result<Vec<_>>>()
})
}
fn mirror_for_primary_records(
primary_base_record_id: u64,
primary_record_count: u64,
mirror: Option<&SequentialMftMirrorChunk>,
record_size: usize,
) -> Option<SequentialMftMirrorChunk> {
let mirror = mirror?;
let primary_end = primary_base_record_id.checked_add(primary_record_count)?;
let mirror_record_count = mirror.bytes.len() / record_size;
let mirror_end = mirror
.base_record_id
.checked_add(u64::try_from(mirror_record_count).ok()?)?;
if primary_base_record_id >= mirror_end || primary_end <= mirror.base_record_id {
return None;
}
Some(mirror.clone())
}
struct LiveNtfsVolume {
handle: HANDLE,
device_path: String,
mft_data_path: String,
}
impl LiveNtfsVolume {
fn open(capabilities: &NtfsVolumeCapabilities) -> Result<Self> {
let device = wide_null(OsStr::new(&capabilities.device_path));
let share_mode =
FILE_SHARE_MODE(FILE_SHARE_READ.0 | FILE_SHARE_WRITE.0 | FILE_SHARE_DELETE.0);
let flags = FILE_FLAGS_AND_ATTRIBUTES(FILE_FLAG_BACKUP_SEMANTICS.0);
let handle = unsafe {
CreateFileW(
PCWSTR(device.as_ptr()),
windows::Win32::Foundation::GENERIC_READ.0,
share_mode,
None,
OPEN_EXISTING,
flags,
None,
)
}
.map_err(|err| volume_open_error(&capabilities.device_path, &err))?;
Ok(Self {
handle,
device_path: capabilities.device_path.clone(),
mft_data_path: capabilities.mft_data_path.clone(),
})
}
fn ntfs_volume_data(&self) -> Result<NTFS_VOLUME_DATA_BUFFER> {
let mut volume_data = NTFS_VOLUME_DATA_BUFFER::default();
let mut bytes_returned = 0_u32;
unsafe {
DeviceIoControl(
self.handle,
FSCTL_GET_NTFS_VOLUME_DATA,
None,
0,
Some((&mut volume_data as *mut NTFS_VOLUME_DATA_BUFFER).cast()),
size_of::<NTFS_VOLUME_DATA_BUFFER>() as u32,
Some(&mut bytes_returned),
None,
)
}
.map_err(|err| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not read NTFS volume data from {}: {}",
self.device_path,
err.message()
))
})?;
Ok(volume_data)
}
fn usn_journal_state(&self) -> Result<ScanCacheUsnJournalState> {
let mut journal_data = USN_JOURNAL_DATA_V0::default();
let mut bytes_returned = 0_u32;
unsafe {
DeviceIoControl(
self.handle,
FSCTL_QUERY_USN_JOURNAL,
None,
0,
Some((&mut journal_data as *mut USN_JOURNAL_DATA_V0).cast()),
size_of::<USN_JOURNAL_DATA_V0>() as u32,
Some(&mut bytes_returned),
None,
)
}
.map_err(|err| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not query USN journal from {}: {}",
self.device_path,
err.message()
))
})?;
usn_journal_data_to_state(&journal_data, &self.device_path)
}
fn read_usn_changes(
&self,
checkpoint: &ScanCacheUsnCheckpoint,
journal_state: &ScanCacheUsnJournalState,
cancellation: &ScanCancellationToken,
) -> Result<Vec<NtfsVolumeIndexUsnChange>> {
let mut start_usn = i64::try_from(checkpoint.next_usn).map_err(|_| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} cannot replay USN range from oversized checkpoint {} on {}",
checkpoint.next_usn, self.device_path
))
})?;
let journal_next_usn = i64::try_from(journal_state.next_usn).map_err(|_| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} cannot replay USN range to oversized checkpoint {} on {}",
journal_state.next_usn, self.device_path
))
})?;
let mut changes = Vec::new();
while start_usn < journal_next_usn {
check_not_cancelled(cancellation)?;
let mut input = READ_USN_JOURNAL_DATA_V0 {
StartUsn: start_usn,
ReasonMask: u32::MAX,
ReturnOnlyOnClose: 0,
Timeout: 0,
BytesToWaitFor: 0,
UsnJournalID: checkpoint.journal_id,
};
let mut output = vec![0_u8; USN_JOURNAL_READ_BUFFER_BYTES];
let mut bytes_returned = 0_u32;
let result = unsafe {
DeviceIoControl(
self.handle,
FSCTL_READ_USN_JOURNAL,
Some((&mut input as *mut READ_USN_JOURNAL_DATA_V0).cast()),
size_of::<READ_USN_JOURNAL_DATA_V0>() as u32,
Some(output.as_mut_ptr().cast()),
output.len() as u32,
Some(&mut bytes_returned),
None,
)
};
match result {
Ok(()) => {}
Err(err) if windows_error_matches(&err, ERROR_HANDLE_EOF) => {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not replay USN range {}..{} from {}; journal reached EOF before current state",
checkpoint.next_usn, journal_state.next_usn, self.device_path
)));
}
Err(err) => {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not replay USN range {}..{} from {}: {}",
checkpoint.next_usn,
journal_state.next_usn,
self.device_path,
err.message()
)));
}
}
let returned = usize::try_from(bytes_returned).unwrap_or(0);
let (next_start_usn, mut chunk_changes) =
parse_usn_journal_read_buffer(&output[..returned], &self.device_path)?;
if next_start_usn <= u64::try_from(start_usn).unwrap_or(0) {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not replay USN range {}..{} from {}; journal cursor did not advance",
checkpoint.next_usn, journal_state.next_usn, self.device_path
)));
}
changes.append(&mut chunk_changes);
if changes.len() > MAX_NTFS_VOLUME_INDEX_USN_REPLAY_RECORDS {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} skipped persistent volume-index replay after reading more than {} USN records from {}",
MAX_NTFS_VOLUME_INDEX_USN_REPLAY_RECORDS, self.device_path
)));
}
start_usn = i64::try_from(next_start_usn).map_err(|_| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} read oversized next USN {next_start_usn} from {}",
self.device_path
))
})?;
}
Ok(changes)
}
fn read_mft_records(
&self,
volume_data: &NTFS_VOLUME_DATA_BUFFER,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<ParsedNtfsRecords> {
let sequential_source = SequentialMftDataSource { volume: self };
let fsctl_source = FsctlRecordMftSource { volume: self };
read_mft_records_from_sources(
&[&sequential_source, &fsctl_source],
volume_data,
cancellation,
monitor,
)
}
fn open_mft_data_stream(&self) -> Result<LiveNtfsMetadataFile> {
let path = wide_null(OsStr::new(&self.mft_data_path));
let share_mode =
FILE_SHARE_MODE(FILE_SHARE_READ.0 | FILE_SHARE_WRITE.0 | FILE_SHARE_DELETE.0);
let flags =
FILE_FLAGS_AND_ATTRIBUTES(FILE_FLAG_OPEN_REPARSE_POINT.0 | FILE_FLAG_SEQUENTIAL_SCAN.0);
let desired_access = FILE_READ_ATTRIBUTES.0 | SYNCHRONIZE.0;
let handle = unsafe {
CreateFileW(
PCWSTR(path.as_ptr()),
desired_access,
share_mode,
None,
OPEN_EXISTING,
flags,
None,
)
}
.map_err(|err| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} could not open {} read-only: {}",
self.mft_data_path,
err.message()
))
})?;
Ok(LiveNtfsMetadataFile { handle })
}
fn mft_extents(
&self,
mft_data: &LiveNtfsMetadataFile,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<Vec<MftExtent>> {
let mut input = STARTING_VCN_INPUT_BUFFER { StartingVcn: 0 };
let mut output = vec![
0_u8;
offset_of!(RETRIEVAL_POINTERS_BUFFER, Extents)
+ (32 * size_of::<RETRIEVAL_POINTERS_BUFFER_0>())
];
loop {
check_mft_build_progress(cancellation, monitor)?;
let mut bytes_returned = 0_u32;
let result = unsafe {
DeviceIoControl(
mft_data.handle,
FSCTL_GET_RETRIEVAL_POINTERS,
Some((&mut input as *mut STARTING_VCN_INPUT_BUFFER).cast()),
size_of::<STARTING_VCN_INPUT_BUFFER>() as u32,
Some(output.as_mut_ptr().cast()),
output.len() as u32,
Some(&mut bytes_returned),
None,
)
};
match result {
Ok(()) => {
let returned = usize::try_from(bytes_returned).unwrap_or(0);
return parse_retrieval_pointer_extents(&output[..returned]);
}
Err(err)
if windows_error_matches(&err, ERROR_MORE_DATA)
&& output.len() < MAX_RETRIEVAL_POINTER_BUFFER_BYTES =>
{
let next_len = output
.len()
.saturating_mul(2)
.min(MAX_RETRIEVAL_POINTER_BUFFER_BYTES);
output.resize(next_len, 0);
}
Err(err) => {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} could not read $MFT retrieval pointers from {}: {}",
self.mft_data_path,
err.message()
)));
}
}
}
}
fn read_volume_bytes(&self, offset: u64, len: usize) -> Result<Vec<u8>> {
let offset = i64::try_from(offset).map_err(|_| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} volume offset overflowed"
))
})?;
let mut buffer = vec![0_u8; len];
let mut bytes_read = 0_u32;
unsafe {
SetFilePointerEx(self.handle, offset, None, FILE_BEGIN).map_err(|err| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} could not seek {} to byte {offset}: {}",
self.device_path,
err.message()
))
})?;
ReadFile(
self.handle,
Some(&mut buffer),
Some(&mut bytes_read),
None,
)
.map_err(|err| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} could not read {} at byte {offset}: {}",
self.device_path,
err.message()
))
})?;
}
buffer.truncate(usize::try_from(bytes_read).unwrap_or(0));
Ok(buffer)
}
}
impl Drop for LiveNtfsVolume {
fn drop(&mut self) {
unsafe {
let _ = CloseHandle(self.handle);
}
}
}
fn usn_journal_data_to_state(
journal_data: &USN_JOURNAL_DATA_V0,
device_path: &str,
) -> Result<ScanCacheUsnJournalState> {
Ok(ScanCacheUsnJournalState {
journal_id: journal_data.UsnJournalID,
first_usn: nonnegative_usn(journal_data.FirstUsn, "first", device_path)?,
next_usn: nonnegative_usn(journal_data.NextUsn, "next", device_path)?,
})
}
fn parse_usn_journal_read_buffer(
raw: &[u8],
device_path: &str,
) -> Result<(u64, Vec<NtfsVolumeIndexUsnChange>)> {
if raw.len() < size_of::<i64>() {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} read truncated USN journal header from {device_path}"
)));
}
let next_start_usn = read_i64_le(raw, 0)
.and_then(|value| u64::try_from(value).ok())
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} read invalid next USN from {device_path}"
))
})?;
let mut changes = Vec::new();
let mut offset = size_of::<i64>();
while offset < raw.len() {
let Some(record_length) = read_u32_le(raw, offset).map(|value| value as usize) else {
return Err(usn_record_parse_error(device_path, "missing record length"));
};
let record_header_len = offset_of!(USN_RECORD_V2, FileName);
if record_length < record_header_len || offset.saturating_add(record_length) > raw.len() {
return Err(usn_record_parse_error(device_path, "invalid record length"));
}
let major_version = read_u16_le(raw, offset + offset_of!(USN_RECORD_V2, MajorVersion))
.ok_or_else(|| usn_record_parse_error(device_path, "missing major version"))?;
if major_version != 2 {
return Err(usn_record_parse_error(
device_path,
"unsupported USN record version",
));
}
let file_reference =
read_u64_le(raw, offset + offset_of!(USN_RECORD_V2, FileReferenceNumber))
.map(file_reference_from_number)
.ok_or_else(|| usn_record_parse_error(device_path, "missing file reference"))?;
let parent_reference = read_u64_le(
raw,
offset + offset_of!(USN_RECORD_V2, ParentFileReferenceNumber),
)
.map(file_reference_from_number)
.ok_or_else(|| usn_record_parse_error(device_path, "missing parent reference"))?;
let usn = read_i64_le(raw, offset + offset_of!(USN_RECORD_V2, Usn))
.and_then(|value| u64::try_from(value).ok())
.ok_or_else(|| usn_record_parse_error(device_path, "invalid record USN"))?;
changes.push(NtfsVolumeIndexUsnChange {
file_reference,
parent_reference,
usn,
});
offset += record_length;
}
Ok((next_start_usn, changes))
}
fn usn_record_parse_error(device_path: &str, reason: &str) -> RebeccaError {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} read malformed USN journal data from {device_path}: {reason}"
))
}
fn read_u16_le(raw: &[u8], offset: usize) -> Option<u16> {
let bytes = raw.get(offset..offset + size_of::<u16>())?;
Some(u16::from_le_bytes(bytes.try_into().ok()?))
}
fn read_u32_le(raw: &[u8], offset: usize) -> Option<u32> {
let bytes = raw.get(offset..offset + size_of::<u32>())?;
Some(u32::from_le_bytes(bytes.try_into().ok()?))
}
fn read_u64_le(raw: &[u8], offset: usize) -> Option<u64> {
let bytes = raw.get(offset..offset + size_of::<u64>())?;
Some(u64::from_le_bytes(bytes.try_into().ok()?))
}
fn read_i64_le(raw: &[u8], offset: usize) -> Option<i64> {
let bytes = raw.get(offset..offset + size_of::<i64>())?;
Some(i64::from_le_bytes(bytes.try_into().ok()?))
}
fn nonnegative_usn(value: i64, label: &str, device_path: &str) -> Result<u64> {
u64::try_from(value).map_err(|_| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} read negative {label} USN {value} from {device_path}"
))
})
}
struct LiveNtfsMetadataFile {
handle: HANDLE,
}
impl Drop for LiveNtfsMetadataFile {
fn drop(&mut self) {
unsafe {
let _ = CloseHandle(self.handle);
}
}
}
struct LiveNtfsIndexStreamSource<'a, 'observer> {
volume: &'a LiveNtfsVolume,
cancellation: &'a ScanCancellationToken,
monitor: &'a NtfsMftBuildMonitor<'observer>,
}
impl NtfsStreamSource for LiveNtfsIndexStreamSource<'_, '_> {
type Error = RebeccaError;
fn read_bytes_at(
&mut self,
volume_offset: u64,
len: usize,
) -> std::result::Result<Vec<u8>, Self::Error> {
check_mft_build_progress(self.cancellation, self.monitor)?;
let bytes = self.volume.read_volume_bytes(volume_offset, len)?;
self.monitor
.add_metric_usize(NtfsMftBuildMetric::StreamReadBytes, bytes.len());
self.monitor.add_metric(NtfsMftBuildMetric::StreamReads, 1);
Ok(bytes)
}
fn should_continue_stream_reads(&self) -> bool {
!self.cancellation.is_cancelled() && !self.monitor.is_timed_out()
}
}
struct SequentialMftDataSource<'a> {
volume: &'a LiveNtfsVolume,
}
impl MftRecordSource for SequentialMftDataSource<'_> {
fn label(&self) -> &'static str {
SEQUENTIAL_MFT_SOURCE_LABEL
}
fn read_records(
&self,
volume_data: &NTFS_VOLUME_DATA_BUFFER,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<ParsedNtfsRecords> {
let geometry = NtfsRecordGeometry::from_volume_data(&self.volume.device_path, volume_data)?;
let mft_data = monitor.measure_checked(
NtfsMftBuildStage::SequentialOpenMftData,
cancellation,
|| self.volume.open_mft_data_stream(),
)?;
let extents = monitor.measure_checked(
NtfsMftBuildStage::SequentialReadRetrievalPointers,
cancellation,
|| self.volume.mft_extents(&mft_data, cancellation, monitor),
)?;
let reader = MftRecordReader::new(geometry.record_size, geometry.sector_size);
let mut records = Vec::new();
let mut caveats = Vec::new();
let mut parse_errors = MftParseErrorCaveats::default();
let mft_mirror =
match self.read_mft_mirror_records(volume_data, geometry, cancellation, monitor) {
Ok(mirror) => mirror,
Err(err) => {
caveats.push(mft_mirror_read_failed_caveat(&err));
None
}
};
let mut context = SequentialMftReadContext {
geometry,
reader: &reader,
cancellation,
monitor,
records: &mut records,
parse_errors: &mut parse_errors,
mft_mirror: mft_mirror.as_ref(),
parse_chunks: Vec::with_capacity(sequential_mft_parse_window_chunks()),
parse_window_chunks: sequential_mft_parse_window_chunks(),
};
for extent in extents {
self.read_extent_records(extent, &mut context)?;
}
context.flush_parse_chunks()?;
parse_errors.append_to(&mut caveats);
Ok(ParsedNtfsRecords {
source_label: self.label(),
records,
caveats,
})
}
}
struct SequentialMftReadContext<'a, 'observer> {
geometry: NtfsRecordGeometry,
reader: &'a MftRecordReader,
cancellation: &'a ScanCancellationToken,
monitor: &'a NtfsMftBuildMonitor<'observer>,
records: &'a mut Vec<NtfsParsedRecord>,
parse_errors: &'a mut MftParseErrorCaveats,
mft_mirror: Option<&'a SequentialMftMirrorChunk>,
parse_chunks: Vec<SequentialMftChunk>,
parse_window_chunks: usize,
}
impl SequentialMftReadContext<'_, '_> {
fn push_parse_chunk(&mut self, chunk: SequentialMftChunk) -> Result<()> {
self.parse_chunks.push(chunk);
if self.parse_chunks.len() >= self.parse_window_chunks {
self.flush_parse_chunks()?;
}
Ok(())
}
fn flush_parse_chunks(&mut self) -> Result<()> {
if self.parse_chunks.is_empty() {
return Ok(());
}
let chunks = std::mem::replace(
&mut self.parse_chunks,
Vec::with_capacity(self.parse_window_chunks),
);
let batches = self.monitor.measure_checked(
NtfsMftBuildStage::SequentialParseRecords,
self.cancellation,
|| parse_sequential_mft_chunks(self.reader, self.cancellation, &chunks),
)?;
for batch in batches {
self.records.extend(batch.records);
for err in batch.errors {
self.parse_errors.record(err.record_id, err.error);
}
}
Ok(())
}
}
impl SequentialMftDataSource<'_> {
fn read_mft_mirror_records(
&self,
volume_data: &NTFS_VOLUME_DATA_BUFFER,
geometry: NtfsRecordGeometry,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<Option<SequentialMftMirrorChunk>> {
let Some(plan) = mft_mirror_read_plan(&self.volume.device_path, volume_data, geometry)?
else {
return Ok(None);
};
let bytes = monitor.measure_checked(
NtfsMftBuildStage::SequentialReadMftMirror,
cancellation,
|| {
self.volume
.read_volume_bytes(plan.volume_offset, plan.byte_len)
},
)?;
if bytes.len() != plan.byte_len {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} read only {} of {} requested $MFTMirr bytes from {}",
bytes.len(),
plan.byte_len,
self.volume.device_path
)));
}
monitor.add_metric_usize(NtfsMftBuildMetric::MftMirrorReadBytes, bytes.len());
monitor.add_metric(NtfsMftBuildMetric::MftMirrorReadChunks, 1);
Ok(Some(SequentialMftMirrorChunk {
base_record_id: plan.base_record_id,
bytes,
}))
}
fn read_extent_records(
&self,
extent: MftExtent,
context: &mut SequentialMftReadContext<'_, '_>,
) -> Result<()> {
let geometry = context.geometry;
let extent_stream_offset = extent
.starting_vcn
.checked_mul(geometry.bytes_per_cluster)
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} $MFT extent stream offset overflowed"
))
})?;
if !extent_stream_offset.is_multiple_of(geometry.record_size as u64) {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} $MFT extent is not file-record aligned"
)));
}
let mut next_record_id = extent_stream_offset / geometry.record_size as u64;
let mut volume_offset = extent
.lcn
.checked_mul(geometry.bytes_per_cluster)
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} volume offset overflowed"
))
})?;
let mut bytes_remaining = extent
.cluster_count
.checked_mul(geometry.bytes_per_cluster)
.ok_or_else(|| {
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} extent length overflowed"
))
})?;
while bytes_remaining > 0 && next_record_id < geometry.max_record_count {
check_mft_build_progress(context.cancellation, context.monitor)?;
let records_remaining = geometry.max_record_count.saturating_sub(next_record_id);
let read_len =
next_mft_chunk_len(bytes_remaining, records_remaining, geometry.record_size);
if read_len == 0 {
break;
}
let bytes = context.monitor.measure_checked(
NtfsMftBuildStage::SequentialReadMftBytes,
context.cancellation,
|| self.volume.read_volume_bytes(volume_offset, read_len),
)?;
if bytes.len() != read_len {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {SEQUENTIAL_MFT_SOURCE_LABEL} read only {} of {read_len} requested bytes from {}",
bytes.len(),
self.volume.device_path
)));
}
context
.monitor
.add_metric_usize(NtfsMftBuildMetric::SequentialMftReadBytes, bytes.len());
context
.monitor
.add_metric(NtfsMftBuildMetric::SequentialMftReadChunks, 1);
let read_len = read_len as u64;
let records_read = read_len / geometry.record_size as u64;
let mirror = mirror_for_primary_records(
next_record_id,
records_read,
context.mft_mirror,
geometry.record_size,
);
context.push_parse_chunk(SequentialMftChunk {
base_record_id: next_record_id,
bytes,
mirror,
})?;
next_record_id = next_record_id.saturating_add(records_read);
volume_offset = volume_offset.saturating_add(read_len);
bytes_remaining = bytes_remaining.saturating_sub(read_len);
}
Ok(())
}
}
fn mft_mirror_read_failed_caveat(err: &RebeccaError) -> ParseCaveat {
ParseCaveat::new(
MFT_MIRROR_READ_FAILED_CAVEAT_CODE,
format!(
"sequential NTFS/MFT parsing could not read bounded $MFTMirr recovery bytes; primary $MFT records remain authoritative: {err}"
),
)
}
struct FsctlRecordMftSource<'a> {
volume: &'a LiveNtfsVolume,
}
impl MftRecordSource for FsctlRecordMftSource<'_> {
fn label(&self) -> &'static str {
FSCTL_RECORD_SOURCE_LABEL
}
fn read_records(
&self,
volume_data: &NTFS_VOLUME_DATA_BUFFER,
cancellation: &ScanCancellationToken,
monitor: &NtfsMftBuildMonitor,
) -> Result<ParsedNtfsRecords> {
let geometry = NtfsRecordGeometry::from_volume_data(&self.volume.device_path, volume_data)?;
let mut records = Vec::new();
let mut caveats = Vec::new();
let mut parse_errors = MftParseErrorCaveats::default();
let mut requested_record = 0_u64;
monitor.measure_checked(
NtfsMftBuildStage::FsctlReadParseRecords,
cancellation,
|| {
while requested_record < geometry.max_record_count {
if requested_record.is_multiple_of(256) {
check_mft_build_progress(cancellation, monitor)?;
}
monitor.add_metric(NtfsMftBuildMetric::FsctlRecordAttempts, 1);
match self
.volume
.read_file_record(requested_record, geometry.record_size)
{
Ok(Some((record_id, raw_record))) => {
monitor.add_metric(NtfsMftBuildMetric::FsctlRecordSuccesses, 1);
let parsed_record_id = low_file_reference_number(record_id);
match NtfsParsedRecord::parse_fsctl_file_record(
parsed_record_id,
&raw_record,
geometry.sector_size,
) {
Ok(record) => records.push(record),
Err(err) => parse_errors.record(parsed_record_id, err),
}
requested_record =
parsed_record_id.max(requested_record).saturating_add(1);
}
Ok(None) => break,
Err(err) if windows_error_matches(&err, ERROR_HANDLE_EOF) => break,
Err(err) if windows_error_matches(&err, ERROR_INVALID_PARAMETER) => {
requested_record = requested_record.saturating_add(1);
}
Err(err) => {
return Err(RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} could not read MFT record {requested_record} from {}: {}",
self.volume.device_path,
err.message()
)));
}
}
}
Ok(())
},
)?;
parse_errors.append_to(&mut caveats);
Ok(ParsedNtfsRecords {
source_label: self.label(),
records,
caveats,
})
}
}
impl LiveNtfsVolume {
fn read_file_record(
&self,
record_number: u64,
record_size: usize,
) -> std::result::Result<Option<(u64, Vec<u8>)>, WindowsError> {
let mut input = NTFS_FILE_RECORD_INPUT_BUFFER {
FileReferenceNumber: record_number as i64,
};
let output_size = NTFS_FILE_RECORD_OUTPUT_HEADER_BYTES + record_size;
let mut output = vec![0_u8; output_size];
let mut bytes_returned = 0_u32;
unsafe {
DeviceIoControl(
self.handle,
FSCTL_GET_NTFS_FILE_RECORD,
Some((&mut input as *mut NTFS_FILE_RECORD_INPUT_BUFFER).cast()),
size_of::<NTFS_FILE_RECORD_INPUT_BUFFER>() as u32,
Some(output.as_mut_ptr().cast()),
output.len() as u32,
Some(&mut bytes_returned),
None,
)
}?;
if bytes_returned == 0 {
return Ok(None);
}
let output_header = unsafe {
ptr::read_unaligned(output.as_ptr().cast::<NTFS_FILE_RECORD_OUTPUT_BUFFER>())
};
let record_length = usize::try_from(output_header.FileRecordLength).unwrap_or(0);
if record_length == 0 {
return Ok(None);
}
let record_offset = NTFS_FILE_RECORD_OUTPUT_HEADER_BYTES;
let available = output.len().saturating_sub(record_offset);
let record_length = record_length.min(available);
Ok(Some((
output_header.FileReferenceNumber as u64,
output[record_offset..record_offset + record_length].to_vec(),
)))
}
}
fn volume_open_error(device_path: &str, err: &WindowsError) -> RebeccaError {
let reason = if windows_error_matches(err, ERROR_ACCESS_DENIED) {
"permission denied while opening the volume; run an elevated shell or use a safe fallback backend"
} else {
"could not open the live volume read-only"
};
RebeccaError::PlatformUnavailable(format!(
"{EXPERIMENTAL_NTFS_MFT_BACKEND_LABEL} {reason} for {device_path}: {}",
err.message()
))
}
fn low_file_reference_number(file_reference_number: u64) -> u64 {
file_reference_number & FILE_REFERENCE_LOW_MASK
}
fn file_reference_sequence_number(file_reference_number: u64) -> u16 {
((file_reference_number >> 48) & 0xFFFF) as u16
}
fn file_reference_from_number(file_reference_number: u64) -> NtfsFileReference {
NtfsFileReference::known(
low_file_reference_number(file_reference_number),
file_reference_sequence_number(file_reference_number),
)
}
fn file_reference_number(reference: NtfsFileReference) -> u64 {
let record_id = reference.record_id & FILE_REFERENCE_LOW_MASK;
reference.sequence_number.map_or(record_id, |sequence| {
record_id | (u64::from(sequence) << 48)
})
}
fn root_metadata(path: &Path) -> Result<std::fs::Metadata> {
std::fs::symlink_metadata(path).map_err(|err| {
RebeccaError::ScanFailed(ScanFailure::from_io(
path,
ScanFailurePhase::RootMetadata,
&err,
))
})
}
fn wide_null(value: &OsStr) -> Vec<u16> {
value.encode_wide().chain(std::iter::once(0)).collect()
}
fn wide_buffer_to_string(buffer: &[u16]) -> String {
let len = buffer
.iter()
.position(|character| *character == 0)
.unwrap_or(buffer.len());
OsString::from_wide(&buffer[..len])
.to_string_lossy()
.into_owned()
}
fn windows_error_matches(err: &WindowsError, code: WIN32_ERROR) -> bool {
err.code() == HRESULT::from_win32(code.0)
}
#[cfg(test)]
mod tests {
use std::collections::{BTreeMap, BTreeSet};
use std::io;
use std::mem::offset_of;
use std::time::{Duration, UNIX_EPOCH};
use rebecca_ntfs::{
AttributeType, FileNameNamespace, MftIndex, MftRecordReader, NtfsAttributeStream,
NtfsDataRun, NtfsDirectoryIndex, NtfsFileName, NtfsFileReference, NtfsIndexEntry,
NtfsParsedRecord, NtfsStreamSource, ParseCaveat, PhysicalMetricsAccumulator,
};
use super::{
CachedNtfsVolumeIndex, MAX_MFT_ESTIMATE_CAVEAT_SAMPLES_PER_CODE,
MAX_MFT_PARSE_ERROR_CAVEAT_SAMPLES, MFT_BUILD_TIMING_CAVEAT_CODE, MFT_CAVEAT_SUMMARY_CODE,
MFT_INDEX_ALLOCATION_BUDGET_EXHAUSTED_CAVEAT_CODE, MFT_PERSISTENT_CACHE_MISS_CAVEAT_CODE,
MFT_PERSISTENT_CACHE_WRITE_SKIPPED_CAVEAT_CODE, MftExtent, MftParseErrorCaveats,
MftRecordSource, NTFS_FILE_RECORD_OUTPUT_HEADER_BYTES, NTFS_VOLUME_DATA_BUFFER,
NTFS_VOLUME_INDEX_CACHE_DIR, NTFS_VOLUME_INDEX_PAYLOAD_VERSION, NtfsMftBuildMetric,
NtfsMftBuildMonitor, NtfsMftBuildMonitorEvent, NtfsMftBuildStage, NtfsRecordGeometry,
NtfsVolumeCapabilities, NtfsVolumeIndexCacheKey, NtfsVolumeIndexCacheManifest,
NtfsVolumeIndexFingerprint, NtfsVolumeIndexManifestLookup, NtfsVolumeIndexManifestMiss,
NtfsVolumeIndexManifestStore, NtfsVolumeIndexPayload, NtfsVolumeIndexPayloadLookup,
NtfsVolumeIndexPayloadMiss, NtfsVolumeIndexPayloadRef, NtfsVolumeIndexUsnChange,
ParsedNtfsRecords, PersistentIndexLoad, RETRIEVAL_POINTERS_BUFFER,
RETRIEVAL_POINTERS_BUFFER_0, SEQUENTIAL_MFT_CHUNK_BYTES, SEQUENTIAL_MFT_SOURCE_LABEL,
ScanBackendEvidence, ScanCancellationToken, SequentialMftChunk, SequentialMftMirrorChunk,
TargetedMftRecordResolver, TargetedMftTraversal, TargetedMftTraversalLimits, USN_RECORD_V2,
VolumePaths, build_mft_index_from_records, cache_checksum, check_mft_build_progress,
collect_mft_disk_map_entry, file_reference_from_number, file_reference_number,
is_transient_persistent_cache_caveat, low_file_reference_number, mft_mirror_read_plan,
mirror_for_primary_records, next_mft_chunk_len, parse_retrieval_pointer_extents,
parse_sequential_mft_chunks, parse_usn_journal_read_buffer, persistent_cache_miss_caveat,
persistent_cache_write_skipped_caveat, read_mft_records_from_sources,
stable_ntfs_volume_index_checkpoint, validate_ntfs_volume_index_replay,
validate_ntfs_volume_index_replay_range, with_bounded_mft_caveats,
};
use crate::disk_map::{
DiskMapBackendOptions, DiskMapEntryKind, DiskMapGroup, DiskMapGroupKind, DiskMapSortField,
DiskMapTopEntries,
};
use crate::error::{RebeccaError, Result};
use crate::plan::EstimateProvenance;
use crate::scan::ScanReport;
use crate::scan::backend::{MeasuredScan, ScanBackendKind};
use crate::scan_cache::{ScanCacheUsnCheckpoint, ScanCacheUsnJournalState};
#[test]
fn volume_paths_support_drive_absolute_paths() {
let paths = VolumePaths::from_path(std::path::Path::new("C:\\Temp\\Cache")).unwrap();
assert_eq!(paths.root_path, std::path::PathBuf::from("C:\\"));
assert_eq!(paths.device_path, "\\\\.\\C:");
assert_eq!(paths.mft_data_path, "\\\\?\\C:\\$MFT::$DATA");
}
#[test]
fn volume_paths_reject_relative_paths() {
let err = VolumePaths::from_path(std::path::Path::new("Temp\\Cache")).unwrap_err();
assert!(err.to_string().contains("absolute local path"));
}
#[test]
fn low_file_reference_masks_sequence_bits() {
assert_eq!(low_file_reference_number(0x0001_0000_0000_002A), 42);
}
#[test]
fn file_reference_roundtrips_sequence_bits_for_targeted_fsctl() {
let reference = file_reference_from_number(0x0003_0000_004B_DD21);
assert_eq!(reference, NtfsFileReference::known(0x4B_DD21, 3));
assert_eq!(file_reference_number(reference), 0x0003_0000_004B_DD21);
assert_eq!(
file_reference_number(NtfsFileReference::unknown_sequence(42)),
42
);
}
#[test]
fn ntfs_volume_cache_key_uses_device_path_and_serial() {
let capabilities = ntfs_volume_capabilities("\\\\.\\C:", 7);
assert_eq!(
capabilities.cache_key(),
NtfsVolumeIndexCacheKey::new("\\\\.\\C:", 7)
);
assert_ne!(
capabilities.cache_key(),
NtfsVolumeIndexCacheKey::new("\\\\.\\D:", 7)
);
assert_ne!(
capabilities.cache_key(),
NtfsVolumeIndexCacheKey::new("\\\\.\\C:", 8)
);
}
#[test]
fn ntfs_volume_index_fingerprint_generation_is_stable() {
let capabilities = ntfs_volume_capabilities("\\\\.\\C:", 7);
let mut volume_data = ntfs_volume_data(1024, 512, 4096, 8192);
volume_data.MftStartLcn = 12;
volume_data.Mft2StartLcn = 34;
let geometry =
NtfsRecordGeometry::from_volume_data(&capabilities.device_path, &volume_data).unwrap();
let first =
NtfsVolumeIndexFingerprint::from_volume_data(&capabilities, &volume_data, geometry);
let second =
NtfsVolumeIndexFingerprint::from_volume_data(&capabilities, &volume_data, geometry);
assert_eq!(first, second);
assert_eq!(
first.persistent_generation(),
second.persistent_generation()
);
assert!(first.is_reusable_for(&capabilities));
assert!(!first.is_reusable_for(&ntfs_volume_capabilities("\\\\.\\D:", 7)));
}
#[test]
fn ntfs_volume_index_fingerprint_generation_changes_with_reuse_boundary() {
let capabilities = ntfs_volume_capabilities("\\\\.\\C:", 7);
let mut volume_data = ntfs_volume_data(1024, 512, 4096, 8192);
volume_data.MftStartLcn = 12;
volume_data.Mft2StartLcn = 34;
let geometry =
NtfsRecordGeometry::from_volume_data(&capabilities.device_path, &volume_data).unwrap();
let base =
NtfsVolumeIndexFingerprint::from_volume_data(&capabilities, &volume_data, geometry);
let mut geometry_changed = volume_data;
geometry_changed.BytesPerFileRecordSegment = 2048;
geometry_changed.MftValidDataLength = 16_384;
let changed_geometry =
NtfsRecordGeometry::from_volume_data(&capabilities.device_path, &geometry_changed)
.unwrap();
let geometry_fingerprint = NtfsVolumeIndexFingerprint::from_volume_data(
&capabilities,
&geometry_changed,
changed_geometry,
);
let mut location_changed = volume_data;
location_changed.MftStartLcn = 99;
let location_fingerprint = NtfsVolumeIndexFingerprint::from_volume_data(
&capabilities,
&location_changed,
geometry,
);
let mut mirror_changed = volume_data;
mirror_changed.Mft2StartLcn = 100;
let mirror_fingerprint =
NtfsVolumeIndexFingerprint::from_volume_data(&capabilities, &mirror_changed, geometry);
let mut length_changed = volume_data;
length_changed.MftValidDataLength = 16_384;
let length_geometry =
NtfsRecordGeometry::from_volume_data(&capabilities.device_path, &length_changed)
.unwrap();
let length_fingerprint = NtfsVolumeIndexFingerprint::from_volume_data(
&capabilities,
&length_changed,
length_geometry,
);
assert_ne!(
base.persistent_generation(),
geometry_fingerprint.persistent_generation()
);
assert_ne!(
base.persistent_generation(),
location_fingerprint.persistent_generation()
);
assert_ne!(
base.persistent_generation(),
mirror_fingerprint.persistent_generation()
);
assert_ne!(
base.persistent_generation(),
length_fingerprint.persistent_generation()
);
}
#[test]
fn ntfs_volume_index_manifest_round_trips_and_validates() {
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let manifest = NtfsVolumeIndexCacheManifest::new(
fingerprint.clone(),
"targeted-fsctl",
Some(ScanCacheUsnCheckpoint {
journal_id: 12,
next_usn: 345,
}),
);
let raw = serde_json::to_string(&manifest).unwrap();
let parsed: NtfsVolumeIndexCacheManifest = serde_json::from_str(&raw).unwrap();
assert_eq!(parsed, manifest);
assert_eq!(parsed.validation_miss(&fingerprint), None);
assert_eq!(parsed.generation, fingerprint.persistent_generation());
assert_eq!(parsed.source_label, "targeted-fsctl");
assert_eq!(parsed.payload, None);
assert_eq!(
parsed.usn_checkpoint,
Some(ScanCacheUsnCheckpoint {
journal_id: 12,
next_usn: 345,
})
);
}
#[test]
fn ntfs_volume_index_manifest_rejects_future_version_and_stale_fingerprint() {
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let manifest =
NtfsVolumeIndexCacheManifest::new(fingerprint.clone(), "targeted-fsctl", None);
let mut future = manifest.clone();
future.version = 999;
assert_eq!(
future.validation_miss(&fingerprint),
Some(NtfsVolumeIndexManifestMiss::UnsupportedVersion)
);
let stale_fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 2048, 512, 4096, 8192);
assert_eq!(
manifest.validation_miss(&stale_fingerprint),
Some(NtfsVolumeIndexManifestMiss::Stale)
);
}
#[test]
fn ntfs_volume_index_payload_round_trips_and_validates() {
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = fixture_mft_index();
let payload = NtfsVolumeIndexPayload::new(
fingerprint.clone(),
"sequential-mft",
vec![ParseCaveat::new("sample", "sample caveat")],
index.clone(),
);
let raw = serde_json::to_vec(&payload).unwrap();
let payload_ref = NtfsVolumeIndexPayloadRef::new("payload.json".to_string(), &raw);
let parsed: NtfsVolumeIndexPayload = serde_json::from_slice(&raw).unwrap();
assert_eq!(parsed, payload);
assert_eq!(parsed.validation_miss(&fingerprint), None);
assert_eq!(payload_ref.validation_miss(&raw), None);
assert_eq!(parsed.mft_index, index);
assert_eq!(parsed.mft_index.aggregate_subtree(5).bytes, 13);
assert_eq!(parsed.mft_index.aggregate_subtree(5).files, 2);
assert_eq!(payload_ref.version, NTFS_VOLUME_INDEX_PAYLOAD_VERSION);
assert_eq!(payload_ref.byte_len, raw.len() as u64);
assert_eq!(payload_ref.checksum, cache_checksum(&raw));
}
#[test]
fn ntfs_volume_index_payload_rejects_future_version_stale_fingerprint_and_checksum_mismatch() {
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let payload = NtfsVolumeIndexPayload::new(
fingerprint.clone(),
"sequential-mft",
Vec::new(),
fixture_mft_index(),
);
let raw = serde_json::to_vec(&payload).unwrap();
let payload_ref = NtfsVolumeIndexPayloadRef::new("payload.json".to_string(), &raw);
let mut future = payload.clone();
future.version = 999;
assert_eq!(
future.validation_miss(&fingerprint),
Some(NtfsVolumeIndexPayloadMiss::UnsupportedVersion)
);
let stale_fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 2048, 512, 4096, 8192);
assert_eq!(
payload.validation_miss(&stale_fingerprint),
Some(NtfsVolumeIndexPayloadMiss::Stale)
);
let mut corrupted = raw.clone();
corrupted.push(0);
assert_eq!(
payload_ref.validation_miss(&corrupted),
Some(NtfsVolumeIndexPayloadMiss::ChecksumMismatch)
);
}
#[test]
fn ntfs_volume_index_manifest_store_round_trips_under_generation_path() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let manifest =
NtfsVolumeIndexCacheManifest::new(fingerprint.clone(), "targeted-fsctl", None);
assert_eq!(
store.load(&fingerprint),
NtfsVolumeIndexManifestLookup::Miss(NtfsVolumeIndexManifestMiss::Missing)
);
store.store(&manifest).unwrap();
let cache_file = store.cache_file_for(&fingerprint);
assert_eq!(
cache_file.parent().unwrap().file_name().unwrap(),
NTFS_VOLUME_INDEX_CACHE_DIR
);
assert_eq!(
cache_file.file_name().unwrap().to_string_lossy(),
format!("{:016x}.json", fingerprint.persistent_generation())
);
assert_eq!(
store.load(&fingerprint),
NtfsVolumeIndexManifestLookup::Hit(manifest)
);
}
#[test]
fn ntfs_volume_index_manifest_store_prunes_corrupt_and_stale_records() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let cache_file = store.cache_file_for(&fingerprint);
std::fs::create_dir_all(cache_file.parent().unwrap()).unwrap();
std::fs::write(&cache_file, b"not-json").unwrap();
assert_eq!(
store.load(&fingerprint),
NtfsVolumeIndexManifestLookup::Miss(NtfsVolumeIndexManifestMiss::Corrupted)
);
assert!(!cache_file.exists());
let stale_fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 2048, 512, 4096, 8192);
let stale_file = store.cache_file_for(&stale_fingerprint);
let manifest =
NtfsVolumeIndexCacheManifest::new(fingerprint.clone(), "targeted-fsctl", None);
std::fs::create_dir_all(stale_file.parent().unwrap()).unwrap();
std::fs::write(&stale_file, serde_json::to_vec(&manifest).unwrap()).unwrap();
assert_eq!(
store.load(&stale_fingerprint),
NtfsVolumeIndexManifestLookup::Miss(NtfsVolumeIndexManifestMiss::Stale)
);
assert!(!stale_file.exists());
}
#[test]
fn ntfs_volume_index_manifest_store_round_trips_payload_pair() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = fixture_mft_index();
let caveats = vec![ParseCaveat::new("sample", "sample caveat")];
assert_eq!(
store.load_index_payload(&fingerprint),
NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Manifest(
NtfsVolumeIndexManifestMiss::Missing
))
);
store
.store_index_payload(
fingerprint.clone(),
"sequential-mft",
None,
&caveats,
&index,
)
.unwrap();
let manifest_file = store.cache_file_for(&fingerprint);
let payload_file = store.payload_file_for(&fingerprint);
assert_eq!(
payload_file.file_name().unwrap().to_string_lossy(),
format!("{:016x}.index.json", fingerprint.persistent_generation())
);
match store.load_index_payload(&fingerprint) {
NtfsVolumeIndexPayloadLookup::Hit { manifest, payload } => {
let payload_ref = manifest.payload.unwrap();
assert_eq!(
payload_ref.file_name,
store.payload_file_name_for(&fingerprint)
);
assert_eq!(payload_ref.version, NTFS_VOLUME_INDEX_PAYLOAD_VERSION);
assert_eq!(manifest_file, store.cache_file_for(&fingerprint));
assert_eq!(payload.mft_index, index);
assert_eq!(payload.caveats, caveats);
assert_eq!(payload.source_label, "sequential-mft");
}
lookup => panic!("expected payload hit, got {lookup:?}"),
}
}
#[test]
fn ntfs_volume_index_stable_checkpoint_requires_unchanged_usn() {
let before = usn_state(11, 100, 200);
let after = usn_state(11, 100, 200);
assert_eq!(
stable_ntfs_volume_index_checkpoint(Some(&before), Some(&after)),
Some(ScanCacheUsnCheckpoint {
journal_id: 11,
next_usn: 200
})
);
assert_eq!(
stable_ntfs_volume_index_checkpoint(Some(&before), Some(&usn_state(12, 100, 200))),
None
);
assert_eq!(
stable_ntfs_volume_index_checkpoint(Some(&before), Some(&usn_state(11, 100, 201))),
None
);
assert_eq!(
stable_ntfs_volume_index_checkpoint(Some(&before), Some(&usn_state(11, 300, 200))),
None
);
assert_eq!(
stable_ntfs_volume_index_checkpoint(None, Some(&after)),
None
);
}
#[test]
fn ntfs_volume_index_replay_range_validation_accepts_advanced_journal() {
let checkpoint = ScanCacheUsnCheckpoint {
journal_id: 11,
next_usn: 200,
};
assert_eq!(
validate_ntfs_volume_index_replay_range(Some(&checkpoint), &usn_state(11, 100, 200)),
None
);
assert_eq!(
validate_ntfs_volume_index_replay_range(None, &usn_state(11, 100, 200)),
Some(NtfsVolumeIndexPayloadMiss::UsnCheckpointMissing)
);
assert_eq!(
validate_ntfs_volume_index_replay_range(Some(&checkpoint), &usn_state(12, 100, 200)),
Some(NtfsVolumeIndexPayloadMiss::UsnJournalChanged)
);
assert_eq!(
validate_ntfs_volume_index_replay_range(Some(&checkpoint), &usn_state(11, 201, 201)),
Some(NtfsVolumeIndexPayloadMiss::UsnRangeUnavailable)
);
assert_eq!(
validate_ntfs_volume_index_replay_range(Some(&checkpoint), &usn_state(11, 100, 199)),
Some(NtfsVolumeIndexPayloadMiss::UsnRangeUnavailable)
);
assert_eq!(
validate_ntfs_volume_index_replay_range(Some(&checkpoint), &usn_state(11, 100, 201)),
None
);
}
#[test]
fn ntfs_volume_index_manifest_store_loads_fresh_payload_with_usn_checkpoint() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = fixture_mft_index();
let checkpoint = ScanCacheUsnCheckpoint {
journal_id: 11,
next_usn: 200,
};
store
.store_index_payload(
fingerprint.clone(),
"sequential-mft",
Some(checkpoint.clone()),
&[],
&index,
)
.unwrap();
match store.load_replayable_index_payload(&fingerprint, &usn_state(11, 100, 200)) {
NtfsVolumeIndexPayloadLookup::Hit { manifest, payload } => {
assert_eq!(manifest.usn_checkpoint, Some(checkpoint));
assert_eq!(payload.mft_index, index);
}
lookup => panic!("expected fresh payload hit, got {lookup:?}"),
}
}
#[test]
fn ntfs_volume_index_payload_store_filters_transient_cache_diagnostics() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let checkpoint = ScanCacheUsnCheckpoint {
journal_id: 11,
next_usn: 200,
};
let index = CachedNtfsVolumeIndex {
fingerprint: fingerprint.clone(),
mft_index: fixture_mft_index(),
source_label: SEQUENTIAL_MFT_SOURCE_LABEL,
caveats: vec![
persistent_cache_miss_caveat("manifest-missing"),
ParseCaveat::new("durable-caveat", "durable"),
],
backend_evidence: ScanBackendEvidence::default(),
usn_checkpoint: Some(checkpoint),
};
assert!(index.store_persistent_payload(Some(&store)).is_none());
match store.load_index_payload(&fingerprint) {
NtfsVolumeIndexPayloadLookup::Hit { payload, .. } => {
assert_eq!(payload.caveats.len(), 1);
assert_eq!(payload.caveats[0].code, "durable-caveat");
}
lookup => panic!("expected payload hit, got {lookup:?}"),
}
assert!(is_transient_persistent_cache_caveat(
MFT_PERSISTENT_CACHE_MISS_CAVEAT_CODE
));
assert!(is_transient_persistent_cache_caveat(
MFT_PERSISTENT_CACHE_WRITE_SKIPPED_CAVEAT_CODE
));
assert!(!is_transient_persistent_cache_caveat("durable-caveat"));
}
#[test]
fn ntfs_volume_index_manifest_store_rejects_checkpoint_free_payload_as_not_fresh() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = fixture_mft_index();
store
.store_index_payload(fingerprint.clone(), "sequential-mft", None, &[], &index)
.unwrap();
assert_eq!(
store.load_replayable_index_payload(&fingerprint, &usn_state(11, 100, 200)),
NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::UsnCheckpointMissing)
);
assert!(store.cache_file_for(&fingerprint).exists());
assert!(store.payload_file_for(&fingerprint).exists());
}
#[test]
fn ntfs_volume_index_manifest_store_keeps_payload_after_usn_advance_for_replay() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = fixture_mft_index();
let checkpoint = ScanCacheUsnCheckpoint {
journal_id: 11,
next_usn: 200,
};
store
.store_index_payload(
fingerprint.clone(),
"sequential-mft",
Some(checkpoint),
&[],
&index,
)
.unwrap();
assert!(store.cache_file_for(&fingerprint).exists());
assert!(store.payload_file_for(&fingerprint).exists());
assert!(matches!(
store.load_replayable_index_payload(&fingerprint, &usn_state(11, 100, 201)),
NtfsVolumeIndexPayloadLookup::Hit { .. }
));
}
#[test]
fn ntfs_volume_index_usn_read_buffer_parses_v2_records() {
let raw = usn_read_buffer(
230,
vec![
usn_record_v2(
NtfsFileReference::known(9, 9),
NtfsFileReference::known(8, 8),
210,
),
usn_record_v2(
NtfsFileReference::known(10, 10),
NtfsFileReference::known(5, 5),
220,
),
],
);
let (next_usn, changes) = parse_usn_journal_read_buffer(&raw, "\\\\.\\C:").unwrap();
assert_eq!(next_usn, 230);
assert_eq!(
changes,
vec![
NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(9, 9),
parent_reference: NtfsFileReference::known(8, 8),
usn: 210,
},
NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(10, 10),
parent_reference: NtfsFileReference::known(5, 5),
usn: 220,
},
]
);
}
#[test]
fn ntfs_volume_index_usn_read_buffer_rejects_unsupported_or_truncated_records() {
let mut unsupported = usn_record_v2(
NtfsFileReference::known(9, 9),
NtfsFileReference::known(8, 8),
210,
);
unsupported[offset_of!(USN_RECORD_V2, MajorVersion)..][..2]
.copy_from_slice(&3_u16.to_le_bytes());
assert!(
parse_usn_journal_read_buffer(&usn_read_buffer(230, vec![unsupported]), "\\\\.\\C:")
.is_err()
);
let mut truncated = usn_read_buffer(
230,
vec![usn_record_v2(
NtfsFileReference::known(9, 9),
NtfsFileReference::known(8, 8),
210,
)],
);
truncated.pop();
assert!(parse_usn_journal_read_buffer(&truncated, "\\\\.\\C:").is_err());
let mut negative_usn = usn_record_v2(
NtfsFileReference::known(9, 9),
NtfsFileReference::known(8, 8),
210,
);
negative_usn[offset_of!(USN_RECORD_V2, Usn)..][..8]
.copy_from_slice(&(-1_i64).to_le_bytes());
assert!(
parse_usn_journal_read_buffer(&usn_read_buffer(230, vec![negative_usn]), "\\\\.\\C:")
.is_err()
);
}
#[test]
fn ntfs_volume_index_usn_replay_accepts_unrelated_subtree_changes() {
let index = nested_fixture_mft_index();
let changes = vec![NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(30, 30),
parent_reference: NtfsFileReference::known(20, 20),
usn: 210,
}];
assert_eq!(
validate_ntfs_volume_index_replay(&index, 8, false, 200, 220, &changes),
None
);
}
#[test]
fn ntfs_volume_index_usn_replay_rejects_target_and_descendant_changes() {
let index = nested_fixture_mft_index();
assert_eq!(
validate_ntfs_volume_index_replay(
&index,
8,
false,
200,
220,
&[NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(8, 8),
parent_reference: NtfsFileReference::known(5, 5),
usn: 210,
}],
),
Some(NtfsVolumeIndexPayloadMiss::UsnTargetChanged)
);
assert_eq!(
validate_ntfs_volume_index_replay(
&index,
8,
false,
200,
220,
&[NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(9, 9),
parent_reference: NtfsFileReference::known(8, 8),
usn: 210,
}],
),
Some(NtfsVolumeIndexPayloadMiss::UsnTargetChanged)
);
}
#[test]
fn ntfs_volume_index_usn_replay_rejects_unknown_ancestry_and_root_changes() {
let index = nested_fixture_mft_index();
assert_eq!(
validate_ntfs_volume_index_replay(
&index,
8,
false,
200,
220,
&[NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(40, 40),
parent_reference: NtfsFileReference::known(99, 99),
usn: 210,
}],
),
Some(NtfsVolumeIndexPayloadMiss::UsnAncestryUnavailable)
);
assert_eq!(
validate_ntfs_volume_index_replay(
&index,
5,
true,
200,
220,
&[NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(30, 30),
parent_reference: NtfsFileReference::known(20, 20),
usn: 210,
}],
),
Some(NtfsVolumeIndexPayloadMiss::UsnTargetChanged)
);
}
#[test]
fn ntfs_volume_index_usn_replay_rejects_hardlink_candidate_under_target() {
let mut file = parsed_file(30, 20, "outside.bin", 1);
file.names
.push(parsed_file_name(8, "inside.bin", FILE_ATTRIBUTE_NORMAL));
let index = MftIndex::from_parsed_records(vec![
parsed_directory(5, 5, "root"),
parsed_directory(8, 5, "target"),
parsed_directory(20, 5, "outside"),
file,
]);
assert_eq!(
validate_ntfs_volume_index_replay(
&index,
8,
false,
200,
220,
&[NtfsVolumeIndexUsnChange {
file_reference: NtfsFileReference::known(30, 30),
parent_reference: NtfsFileReference::known(20, 20),
usn: 210,
}],
),
Some(NtfsVolumeIndexPayloadMiss::UsnTargetChanged)
);
}
#[test]
fn ntfs_volume_index_payload_loader_treats_manifest_only_record_as_miss() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let manifest =
NtfsVolumeIndexCacheManifest::new(fingerprint.clone(), "targeted-fsctl", None);
store.store(&manifest).unwrap();
assert_eq!(
store.load_index_payload(&fingerprint),
NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::ManifestWithoutPayload)
);
assert!(store.cache_file_for(&fingerprint).exists());
}
#[test]
fn ntfs_volume_index_payload_loader_prunes_missing_and_corrupt_payload_pairs() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = fixture_mft_index();
store
.store_index_payload(fingerprint.clone(), "sequential-mft", None, &[], &index)
.unwrap();
let manifest_file = store.cache_file_for(&fingerprint);
let payload_file = store.payload_file_for(&fingerprint);
std::fs::remove_file(&payload_file).unwrap();
assert_eq!(
store.load_index_payload(&fingerprint),
NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Missing)
);
assert!(!manifest_file.exists());
let raw = b"not-json";
let payload_ref =
NtfsVolumeIndexPayloadRef::new(store.payload_file_name_for(&fingerprint), raw);
let manifest =
NtfsVolumeIndexCacheManifest::new(fingerprint.clone(), "sequential-mft", None)
.with_payload(payload_ref);
store.store(&manifest).unwrap();
std::fs::write(&payload_file, raw).unwrap();
assert_eq!(
store.load_index_payload(&fingerprint),
NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Corrupted)
);
assert!(!manifest_file.exists());
assert!(!payload_file.exists());
}
#[test]
fn ntfs_volume_index_payload_loader_prunes_stale_payload_pair() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let stale_fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 2048, 512, 4096, 8192);
let payload = NtfsVolumeIndexPayload::new(
fingerprint,
"sequential-mft",
Vec::new(),
fixture_mft_index(),
);
let raw = serde_json::to_vec(&payload).unwrap();
let payload_ref =
NtfsVolumeIndexPayloadRef::new(store.payload_file_name_for(&stale_fingerprint), &raw);
let manifest =
NtfsVolumeIndexCacheManifest::new(stale_fingerprint.clone(), "sequential-mft", None)
.with_payload(payload_ref);
let payload_file = store.payload_file_for(&stale_fingerprint);
store.store(&manifest).unwrap();
std::fs::write(&payload_file, raw).unwrap();
assert_eq!(
store.load_index_payload(&stale_fingerprint),
NtfsVolumeIndexPayloadLookup::Miss(NtfsVolumeIndexPayloadMiss::Stale)
);
assert!(!store.cache_file_for(&stale_fingerprint).exists());
assert!(!payload_file.exists());
}
#[test]
fn ntfs_volume_index_cache_can_be_configured_with_manifest_store() {
let temp = tempfile::tempdir().unwrap();
let cache = super::WindowsNtfsMftIndexCache::with_manifest_store(
NtfsVolumeIndexManifestStore::new(temp.path()),
);
assert!(cache.manifest_store.is_some());
assert!(
super::WindowsNtfsMftIndexCache::default()
.manifest_store
.is_none()
);
}
#[test]
fn persistent_cache_caveats_use_stable_reason_labels() {
let miss = persistent_cache_miss_caveat("manifest-missing");
assert_eq!(miss.code, MFT_PERSISTENT_CACHE_MISS_CAVEAT_CODE);
assert!(miss.message.contains("reason=manifest-missing"));
let load = PersistentIndexLoad::miss("manifest-missing");
assert_eq!(
load.backend_evidence.cache_events[0].cache,
"ntfs-volume-index"
);
assert_eq!(load.backend_evidence.cache_events[0].outcome, "miss");
assert_eq!(
load.backend_evidence.cache_events[0].reason.as_deref(),
Some("manifest-missing")
);
let write_skip = persistent_cache_write_skipped_caveat("write-failed");
assert_eq!(
write_skip.code,
MFT_PERSISTENT_CACHE_WRITE_SKIPPED_CAVEAT_CODE
);
assert!(write_skip.message.contains("reason=write-failed"));
}
#[test]
fn ntfs_volume_index_store_reports_skipped_payload_without_checkpoint() {
let temp = tempfile::tempdir().unwrap();
let store = NtfsVolumeIndexManifestStore::new(temp.path().join("cache"));
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = CachedNtfsVolumeIndex {
fingerprint,
mft_index: fixture_mft_index(),
source_label: SEQUENTIAL_MFT_SOURCE_LABEL,
caveats: Vec::new(),
backend_evidence: ScanBackendEvidence::default(),
usn_checkpoint: None,
};
let (caveat, reason) = index.store_persistent_payload(Some(&store)).unwrap();
assert_eq!(caveat.code, MFT_PERSISTENT_CACHE_WRITE_SKIPPED_CAVEAT_CODE);
assert_eq!(reason, "stable-usn-checkpoint-unavailable");
assert!(
caveat
.message
.contains("reason=stable-usn-checkpoint-unavailable")
);
assert!(!store.cache_file_for(&index.fingerprint).exists());
}
#[test]
fn ntfs_volume_index_store_without_manifest_store_stays_quiet() {
let fingerprint = ntfs_volume_fingerprint("\\\\.\\C:", 7, 1024, 512, 4096, 8192);
let index = CachedNtfsVolumeIndex {
fingerprint,
mft_index: fixture_mft_index(),
source_label: SEQUENTIAL_MFT_SOURCE_LABEL,
caveats: Vec::new(),
backend_evidence: ScanBackendEvidence::default(),
usn_checkpoint: None,
};
assert!(index.store_persistent_payload(None).is_none());
}
#[test]
fn ntfs_file_record_output_buffer_uses_wire_header_offset() {
assert_eq!(NTFS_FILE_RECORD_OUTPUT_HEADER_BYTES, 12);
}
#[test]
fn ntfs_record_geometry_accepts_valid_volume_data() {
let volume_data = ntfs_volume_data(1024, 512, 4096, 8192);
let geometry = NtfsRecordGeometry::from_volume_data("\\\\.\\C:", &volume_data).unwrap();
assert_eq!(geometry.record_size, 1024);
assert_eq!(geometry.sector_size, 512);
assert_eq!(geometry.bytes_per_cluster, 4096);
assert_eq!(geometry.max_record_count, 8);
}
#[test]
fn ntfs_record_geometry_rejects_unaligned_records() {
let volume_data = ntfs_volume_data(1000, 512, 4096, 8192);
let err = NtfsRecordGeometry::from_volume_data("\\\\.\\C:", &volume_data).unwrap_err();
assert!(err.to_string().contains("unaligned"));
}
#[test]
fn mft_mirror_read_plan_uses_mft2_lcn_for_bounded_system_records() {
let mut volume_data = ntfs_volume_data(1024, 512, 4096, 8192);
volume_data.Mft2StartLcn = 42;
let geometry = NtfsRecordGeometry::from_volume_data("\\\\.\\C:", &volume_data).unwrap();
let plan = mft_mirror_read_plan("\\\\.\\C:", &volume_data, geometry)
.unwrap()
.unwrap();
assert_eq!(plan.base_record_id, 0);
assert_eq!(plan.record_count, 4);
assert_eq!(plan.volume_offset, 42 * 4096);
assert_eq!(plan.byte_len, 4 * 1024);
}
#[test]
fn mft_mirror_read_plan_skips_absent_or_empty_mirror() {
let volume_data = ntfs_volume_data(1024, 512, 4096, 8192);
let geometry = NtfsRecordGeometry::from_volume_data("\\\\.\\C:", &volume_data).unwrap();
assert!(
mft_mirror_read_plan("\\\\.\\C:", &volume_data, geometry)
.unwrap()
.is_none()
);
}
#[test]
fn mft_mirror_chunk_is_attached_only_to_overlapping_primary_records() {
let mirror = SequentialMftMirrorChunk {
base_record_id: 0,
bytes: vec![0xAB; 4 * 1024],
};
assert!(mirror_for_primary_records(10, 2, Some(&mirror), 1024).is_none());
let attached = mirror_for_primary_records(2, 4, Some(&mirror), 1024).unwrap();
assert_eq!(attached.base_record_id, 0);
assert_eq!(attached.bytes.len(), 4 * 1024);
}
#[test]
fn retrieval_pointer_parser_maps_ordered_extents() {
let buffer = retrieval_pointer_buffer(0, &[(4, 10), (9, 20)]);
let extents = parse_retrieval_pointer_extents(&buffer).unwrap();
assert_eq!(
extents,
vec![
MftExtent {
starting_vcn: 0,
lcn: 10,
cluster_count: 4,
},
MftExtent {
starting_vcn: 4,
lcn: 20,
cluster_count: 5,
},
]
);
}
#[test]
fn retrieval_pointer_parser_rejects_sparse_extents() {
let buffer = retrieval_pointer_buffer(0, &[(4, -1)]);
let err = parse_retrieval_pointer_extents(&buffer).unwrap_err();
assert!(err.to_string().contains("sparse"));
}
#[test]
fn mft_chunk_len_is_bounded_and_record_aligned() {
assert_eq!(next_mft_chunk_len(4097, 10, 1024), 4096);
assert_eq!(
next_mft_chunk_len(SEQUENTIAL_MFT_CHUNK_BYTES as u64 * 2, 100_000, 1024),
SEQUENTIAL_MFT_CHUNK_BYTES
);
assert_eq!(next_mft_chunk_len(512, 10, 1024), 0);
}
#[test]
fn mft_index_builder_expands_live_index_allocation_streams() {
let monitor = test_build_monitor();
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let records = ParsedNtfsRecords {
source_label: "sequential",
records: vec![
parsed_directory_with_index_allocation(5, "large-dir"),
parsed_file(6, 99, "large.bin", 3),
],
caveats: vec![ParseCaveat::new("source-caveat", "source")],
};
let (index, caveats) = build_mft_index_from_records(
records,
test_record_geometry(),
&mut source,
&ScanCancellationToken::new(),
&monitor,
)
.unwrap();
let summary = index.aggregate_subtree(5);
assert_eq!(summary.bytes, 3);
assert!(summary.caveats.iter().any(|caveat| {
caveat.code == "directory-index-parent-map-fallback"
&& caveat.message.contains("large.bin")
}));
assert_eq!(caveats.len(), 1);
assert_eq!(caveats[0].code, "source-caveat");
}
#[test]
fn mft_index_builder_turns_stream_read_failure_into_bounded_caveat() {
let monitor = test_build_monitor();
let mut source = FakeIndexStreamSource::default();
let records = ParsedNtfsRecords {
source_label: "sequential",
records: vec![parsed_directory_with_index_allocation(5, "large-dir")],
caveats: Vec::new(),
};
let (index, caveats) = build_mft_index_from_records(
records,
test_record_geometry(),
&mut source,
&ScanCancellationToken::new(),
&monitor,
)
.unwrap();
let summary = index.aggregate_subtree(5);
assert!(caveats.is_empty());
assert!(summary.caveats.iter().any(|caveat| {
caveat.code == "invalid-index-allocation"
&& caveat.message.contains("stream read failed")
}));
}
#[test]
fn mft_index_builder_preserves_cancellation_during_stream_expansion() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut source = CancellingIndexStreamSource {
cancellation: cancellation.clone(),
};
let records = ParsedNtfsRecords {
source_label: "sequential",
records: vec![parsed_directory_with_index_allocation(5, "large-dir")],
caveats: Vec::new(),
};
let err = build_mft_index_from_records(
records,
test_record_geometry(),
&mut source,
&cancellation,
&monitor,
)
.unwrap_err();
assert!(matches!(err, RebeccaError::OperationCancelled(_)));
}
#[test]
fn mft_index_builder_preserves_full_index_evidence_after_index_allocation_budget_exhaustion() {
let mut source = FakeIndexStreamSource::default()
.with_read_delay(Duration::from_millis(150))
.with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let records = ParsedNtfsRecords {
source_label: "sequential",
records: vec![
parsed_directory_with_index_allocation(5, "large-dir"),
parsed_file(6, 99, "large.bin", 3),
],
caveats: vec![ParseCaveat::new("source-caveat", "source")],
};
let monitor = NtfsMftBuildMonitor::near_timeout_for_test(
Duration::from_millis(200),
Duration::from_millis(100),
);
let (index, caveats) = build_mft_index_from_records(
records,
test_record_geometry(),
&mut source,
&ScanCancellationToken::new(),
&monitor,
)
.unwrap();
let summary = index.aggregate_subtree(5);
assert_eq!(summary.bytes, 3);
assert!(
caveats
.iter()
.any(|caveat| caveat.code == MFT_INDEX_ALLOCATION_BUDGET_EXHAUSTED_CAVEAT_CODE)
);
assert!(caveats.iter().any(|caveat| caveat.code == "source-caveat"));
}
#[test]
fn mft_index_builder_stops_later_stream_expansion_when_source_requests_stop() {
let monitor = test_build_monitor();
let mut source = FakeIndexStreamSource::default()
.with_max_reads(1)
.with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"first.bin",
0,
),
index_allocation_last_entry(),
],
),
)
.with_bytes(
0x81_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(8, 8),
file_reference(7, 7),
"second.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut second_directory = parsed_directory_with_index_allocation(7, "second-dir");
second_directory.attribute_streams[0].data_runs[0].lcn = Some(0x81);
let records = ParsedNtfsRecords {
source_label: "sequential",
records: vec![
parsed_directory_with_index_allocation(5, "first-dir"),
parsed_file(6, 99, "first.bin", 3),
second_directory,
parsed_file(8, 99, "second.bin", 5),
],
caveats: Vec::new(),
};
let (index, caveats) = build_mft_index_from_records(
records,
test_record_geometry(),
&mut source,
&ScanCancellationToken::new(),
&monitor,
)
.unwrap();
assert!(caveats.is_empty());
assert_eq!(source.read_count, 1);
assert_eq!(index.aggregate_subtree(5).bytes, 3);
assert_eq!(index.aggregate_subtree(7).bytes, 0);
}
#[test]
fn mft_index_builder_fails_when_budget_expired_before_index_allocation_resolution() {
let mut source = FakeIndexStreamSource::default();
let records = ParsedNtfsRecords {
source_label: "sequential",
records: vec![parsed_directory_with_index_allocation(5, "large-dir")],
caveats: Vec::new(),
};
let monitor = NtfsMftBuildMonitor::expired_for_test(Duration::from_secs(1));
let err = build_mft_index_from_records(
records,
test_record_geometry(),
&mut source,
&ScanCancellationToken::new(),
&monitor,
)
.unwrap_err();
assert!(matches!(err, RebeccaError::PlatformUnavailable(_)));
}
#[test]
fn targeted_traversal_expands_index_allocation_without_full_index() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "large-dir"))
.with_record(parsed_file(6, 99, "large.bin", 3));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let summary = {
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
traversal
.aggregate_subtree(NtfsFileReference::known(5, 5))
.unwrap()
};
assert_eq!(summary.bytes, 3);
assert_eq!(summary.files, 1);
assert_eq!(summary.directories, 1);
assert_eq!(resolver.reads, vec![5, 6]);
}
#[test]
fn targeted_traversal_caveats_child_sequence_mismatch() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "large-dir"))
.with_record(parsed_file(6, 99, "large.bin", 3));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 99),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let summary = traversal
.aggregate_subtree(NtfsFileReference::known(5, 5))
.unwrap();
assert_eq!(summary.bytes, 0);
assert_eq!(summary.files, 0);
assert!(summary.caveats.iter().any(|caveat| {
caveat.code == "directory-index-child-sequence-mismatch"
&& caveat.message.contains("record 6")
}));
}
#[test]
fn targeted_traversal_stops_at_record_budget() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "large-dir"))
.with_record(parsed_file(6, 99, "large.bin", 3));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits {
max_records: 1,
max_depth: 512,
},
};
let err = traversal
.aggregate_subtree(NtfsFileReference::known(5, 5))
.unwrap_err();
assert!(matches!(err, RebeccaError::PlatformUnavailable(_)));
assert!(err.to_string().contains("record candidate budget"));
}
#[test]
fn targeted_traversal_stale_child_does_not_poison_later_valid_child() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "large-dir"))
.with_record(parsed_file(6, 5, "large.bin", 3));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_entry(
file_reference(6, 99),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let summary = traversal
.aggregate_subtree(NtfsFileReference::known(5, 5))
.unwrap();
assert_eq!(summary.bytes, 3);
assert_eq!(summary.files, 1);
assert!(summary.caveats.iter().any(|caveat| {
caveat.code == "directory-index-child-sequence-mismatch"
&& caveat.message.contains("record 6")
}));
}
#[test]
fn targeted_traversal_caveats_i30_parent_map_fallback() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "large-dir"))
.with_record(parsed_file(6, 99, "large.bin", 3));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let summary = traversal
.aggregate_subtree(NtfsFileReference::known(5, 5))
.unwrap();
assert_eq!(summary.bytes, 3);
assert!(summary.caveats.iter().any(|caveat| {
caveat.code == "directory-index-parent-map-fallback"
&& caveat.message.contains("large.bin")
}));
}
#[test]
fn targeted_traversal_skips_dos_i30_aliases() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 5, "long-name.bin", 10));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"long-name.bin",
0,
),
index_allocation_entry_with_namespace(
file_reference(6, 6),
file_reference(5, 5),
"LONG-N~1.BIN",
0,
2,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let summary = traversal
.aggregate_subtree(NtfsFileReference::known(5, 5))
.unwrap();
assert_eq!(summary.bytes, 10);
assert_eq!(summary.files, 1);
assert!(!summary.caveats.iter().any(|caveat| {
caveat.code == "directory-index-parent-map-fallback"
&& caveat.message.contains("LONG-N~1.BIN")
}));
}
#[test]
fn targeted_disk_map_collects_ranked_entries_without_full_index() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 5, "large.bin", 10))
.with_record(parsed_file(7, 5, "small.bin", 3));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(7, 7),
file_reference(5, 5),
"small.bin",
0,
),
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let map = traversal
.collect_disk_map(
NtfsFileReference::known(5, 5),
std::path::Path::new("C:\\root"),
&disk_map_options(2, None, Vec::new()),
100,
&EstimateProvenance::default(),
)
.unwrap();
assert_eq!(map.metrics.logical_bytes, 13);
assert_eq!(map.metrics.allocated_bytes, Some(13));
assert_eq!(map.metrics.unique_logical_bytes, Some(13));
assert_eq!(map.metrics.unique_allocated_bytes, Some(13));
assert_eq!(map.metrics.files, 2);
assert_eq!(map.metrics.directories, 0);
assert_eq!(map.top_entries.len(), 2);
assert_eq!(
map.top_entries[0].path,
std::path::PathBuf::from("C:\\root\\large.bin")
);
assert_eq!(map.top_entries[0].kind, DiskMapEntryKind::File);
assert_eq!(map.top_entries[0].depth, 1);
assert_eq!(
map.top_entries[1].path,
std::path::PathBuf::from("C:\\root\\small.bin")
);
assert!(map.caveats.is_empty());
}
#[test]
fn full_index_disk_map_preserves_duplicate_paths_with_unique_metrics() {
let mut dir_a = parsed_directory_with_index_allocation(6, "a");
dir_a.names = vec![parsed_file_name(5, "a", FILE_ATTRIBUTE_DIRECTORY)];
let mut dir_b = parsed_directory_with_index_allocation(7, "b");
dir_b.names = vec![parsed_file_name(5, "b", FILE_ATTRIBUTE_DIRECTORY)];
let mut file = parsed_file(8, 6, "left.bin", 10);
file.names.push(parsed_file_name(7, "right.bin", 0));
let index = MftIndex::from_parsed_records(vec![
parsed_directory_with_index_allocation(5, "root"),
dir_a,
dir_b,
file,
]);
let options = disk_map_options(10, None, Vec::new());
let mut top_entries = DiskMapTopEntries::new(
options.top_limit,
options.top_sort,
options.entry_filter.clone(),
);
let mut groups = options.group_collector();
let mut visited_directories = BTreeSet::new();
let mut caveats = Vec::new();
let mut aggregate = PhysicalMetricsAccumulator::default();
let cancellation = ScanCancellationToken::new();
let provenance = EstimateProvenance::default();
for edge in index.child_edges(5).cloned().collect::<Vec<_>>() {
let child = index.get(edge.child.record_id).unwrap().clone();
let child_path = std::path::PathBuf::from("C:\\root").join(edge.name);
let child_aggregate = collect_mft_disk_map_entry(
&index,
std::path::Path::new("C:\\root"),
child_path,
child,
1,
usize::MAX,
&provenance,
&mut visited_directories,
&mut caveats,
&mut top_entries,
&mut groups,
100,
&cancellation,
)
.unwrap();
aggregate.absorb_child(child_aggregate);
}
let metrics = aggregate.into_metrics();
assert_eq!(metrics.logical_bytes, 20);
assert_eq!(metrics.allocated_bytes, Some(20));
assert_eq!(metrics.unique_logical_bytes, 10);
assert_eq!(metrics.unique_allocated_bytes, Some(10));
assert_eq!(metrics.files, 2);
assert_eq!(metrics.directories, 2);
let paths = top_entries
.into_sorted_entries()
.into_iter()
.map(|entry| entry.path)
.collect::<BTreeSet<_>>();
assert!(paths.contains(&std::path::PathBuf::from("C:\\root\\a\\left.bin")));
assert!(paths.contains(&std::path::PathBuf::from("C:\\root\\b\\right.bin")));
}
#[test]
fn targeted_disk_map_collects_requested_groups_without_full_index() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 5, "large.bin", 10))
.with_record(parsed_file(7, 5, "small.txt", 3));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_entry(
file_reference(7, 7),
file_reference(5, 5),
"small.txt",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let map = traversal
.collect_disk_map(
NtfsFileReference::known(5, 5),
std::path::Path::new("C:\\root"),
&disk_map_options(
0,
None,
vec![
DiskMapGroupKind::Extension,
DiskMapGroupKind::Depth,
DiskMapGroupKind::Age,
],
),
100,
&EstimateProvenance::default(),
)
.unwrap();
let groups = map.groups.finish();
assert_group_metrics(&groups, DiskMapGroupKind::Extension, ".bin", 10, 1);
assert_group_metrics(&groups, DiskMapGroupKind::Extension, ".txt", 3, 1);
assert_group_metrics(&groups, DiskMapGroupKind::Depth, "depth-1", 13, 2);
assert_group_metrics(&groups, DiskMapGroupKind::Age, "modified-unknown", 13, 2);
}
#[test]
fn targeted_disk_map_max_depth_limits_entries_not_totals() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 5, "large.bin", 10));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let map = traversal
.collect_disk_map(
NtfsFileReference::known(5, 5),
std::path::Path::new("C:\\root"),
&disk_map_options(10, Some(0), Vec::new()),
100,
&EstimateProvenance::default(),
)
.unwrap();
assert_eq!(map.metrics.logical_bytes, 10);
assert_eq!(map.metrics.files, 1);
assert!(map.top_entries.is_empty());
}
#[test]
fn targeted_disk_map_preserves_duplicate_paths_with_unique_metrics() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 5, "large.bin", 10));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"alias.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let map = traversal
.collect_disk_map(
NtfsFileReference::known(5, 5),
std::path::Path::new("C:\\root"),
&disk_map_options(10, None, Vec::new()),
100,
&EstimateProvenance::default(),
)
.unwrap();
assert_eq!(map.metrics.logical_bytes, 20);
assert_eq!(map.metrics.allocated_bytes, Some(20));
assert_eq!(map.metrics.unique_logical_bytes, Some(10));
assert_eq!(map.metrics.unique_allocated_bytes, Some(10));
assert_eq!(map.metrics.files, 2);
assert_eq!(map.top_entries.len(), 2);
let paths = map
.top_entries
.iter()
.map(|entry| entry.path.clone())
.collect::<BTreeSet<_>>();
assert!(paths.contains(&std::path::PathBuf::from("C:\\root\\large.bin")));
assert!(paths.contains(&std::path::PathBuf::from("C:\\root\\alias.bin")));
assert!(
!map.caveats
.iter()
.any(|caveat| caveat.code == "mft-targeted-record-already-counted")
);
}
#[test]
fn targeted_disk_map_skips_dos_i30_aliases() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 5, "long-name.bin", 10));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"long-name.bin",
0,
),
index_allocation_entry_with_namespace(
file_reference(6, 6),
file_reference(5, 5),
"LONG-N~1.BIN",
0,
2,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let map = traversal
.collect_disk_map(
NtfsFileReference::known(5, 5),
std::path::Path::new("C:\\root"),
&disk_map_options(10, None, Vec::new()),
100,
&EstimateProvenance::default(),
)
.unwrap();
assert_eq!(map.metrics.logical_bytes, 10);
assert_eq!(map.metrics.unique_logical_bytes, Some(10));
assert_eq!(map.metrics.files, 1);
assert_eq!(map.top_entries.len(), 1);
assert_eq!(
map.top_entries[0].path,
std::path::PathBuf::from("C:\\root\\long-name.bin")
);
assert!(!map.caveats.iter().any(|caveat| {
caveat.code == "directory-index-parent-map-fallback"
&& caveat.message.contains("LONG-N~1.BIN")
}));
}
#[test]
fn targeted_disk_map_stale_child_does_not_poison_later_valid_child() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 5, "large.bin", 10));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_entry(
file_reference(6, 99),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let map = traversal
.collect_disk_map(
NtfsFileReference::known(5, 5),
std::path::Path::new("C:\\root"),
&disk_map_options(10, None, Vec::new()),
100,
&EstimateProvenance::default(),
)
.unwrap();
assert_eq!(map.metrics.logical_bytes, 10);
assert_eq!(map.metrics.files, 1);
assert_eq!(map.top_entries.len(), 1);
assert_eq!(
map.top_entries[0].path,
std::path::PathBuf::from("C:\\root\\large.bin")
);
assert!(map.caveats.iter().any(|caveat| {
caveat.code == "directory-index-child-sequence-mismatch"
&& caveat.message.contains("record 6")
}));
}
#[test]
fn targeted_disk_map_caveats_i30_parent_map_fallback() {
let monitor = test_build_monitor();
let cancellation = ScanCancellationToken::new();
let mut resolver = FakeTargetedRecordResolver::default()
.with_record(parsed_directory_with_index_allocation(5, "root"))
.with_record(parsed_file(6, 99, "large.bin", 10));
let mut source = FakeIndexStreamSource::default().with_bytes(
0x80_000,
&index_allocation_record(
0,
vec![
index_allocation_entry(
file_reference(6, 6),
file_reference(5, 5),
"large.bin",
0,
),
index_allocation_last_entry(),
],
),
);
let mut traversal = TargetedMftTraversal {
resolver: &mut resolver,
stream_source: &mut source,
geometry: test_record_geometry(),
cancellation: &cancellation,
monitor: &monitor,
limits: TargetedMftTraversalLimits::default(),
};
let map = traversal
.collect_disk_map(
NtfsFileReference::known(5, 5),
std::path::Path::new("C:\\root"),
&disk_map_options(10, None, Vec::new()),
100,
&EstimateProvenance::default(),
)
.unwrap();
assert_eq!(map.metrics.logical_bytes, 10);
assert_eq!(map.metrics.files, 1);
assert!(map.caveats.iter().any(|caveat| {
caveat.code == "directory-index-parent-map-fallback"
&& caveat.message.contains("large.bin")
}));
}
#[test]
fn mft_build_budget_timeout_is_fallback_capable_platform_error() {
let monitor = NtfsMftBuildMonitor::expired_for_test(Duration::from_secs(1));
let err = check_mft_build_progress(&ScanCancellationToken::new(), &monitor).unwrap_err();
assert!(matches!(err, RebeccaError::PlatformUnavailable(_)));
let message = err.to_string();
assert!(message.contains("timed out after 1s"));
assert!(message.contains("REBECCA_NTFS_MFT_INDEX_TIMEOUT_SECONDS"));
}
#[test]
fn mft_build_timeout_reports_active_stage_and_completed_timings() {
let monitor = NtfsMftBuildMonitor::expired_for_test(Duration::from_secs(1));
monitor
.measure(NtfsMftBuildStage::OpenVolume, || Ok(()))
.unwrap();
monitor.add_metric(NtfsMftBuildMetric::ParsedRecords, 2);
let err = monitor
.measure(NtfsMftBuildStage::SequentialParseRecords, || {
check_mft_build_progress(&ScanCancellationToken::new(), &monitor)
})
.unwrap_err();
let message = err.to_string();
assert!(message.contains("while sequential-parse-records"));
assert!(message.contains("completed_timings=open-volume="));
assert!(message.contains("metrics=parsed-records=2"));
}
#[test]
fn mft_build_timing_caveat_is_opt_in() {
let monitor = NtfsMftBuildMonitor::new(None, true);
monitor
.measure(NtfsMftBuildStage::SequentialReadMftBytes, || Ok(()))
.unwrap();
monitor
.measure(NtfsMftBuildStage::SequentialParseRecords, || Ok(()))
.unwrap();
monitor
.measure(NtfsMftBuildStage::BuildMftIndex, || Ok(()))
.unwrap();
monitor.add_metric(NtfsMftBuildMetric::SequentialMftReadBytes, 4096);
monitor.add_metric(NtfsMftBuildMetric::SequentialMftReadChunks, 1);
monitor.add_metric(NtfsMftBuildMetric::ParsedRecords, 4);
let evidence = monitor.evidence();
let caveat = monitor.timing_caveat().unwrap();
assert!(
evidence
.timings_ms
.contains_key("sequential-read-mft-bytes")
);
assert!(evidence.timings_ms.contains_key("sequential-parse-records"));
assert!(evidence.timings_ms.contains_key("build-mft-index"));
assert_eq!(evidence.counters["parsed-records"], 4);
assert_eq!(evidence.counters["sequential-mft-read-bytes"], 4096);
assert_eq!(evidence.counters["sequential-mft-read-chunks"], 1);
assert_eq!(caveat.code, MFT_BUILD_TIMING_CAVEAT_CODE);
assert!(caveat.message.contains("sequential-read-mft-bytes="));
assert!(caveat.message.contains("sequential-parse-records="));
assert!(caveat.message.contains("build-mft-index="));
assert!(caveat.message.contains("metrics=parsed-records=4"));
assert!(caveat.message.contains("sequential-mft-read-bytes=4096"));
assert!(caveat.message.contains("sequential-mft-read-chunks=1"));
}
#[test]
fn mft_build_monitor_observer_receives_stage_and_metric_events() {
let mut events = Vec::new();
{
let mut observer = |event| {
events.push(match event {
NtfsMftBuildMonitorEvent::StageStarted { stage } => {
format!("started:{stage}")
}
NtfsMftBuildMonitorEvent::StageFinished { stage } => {
format!("finished:{stage}")
}
NtfsMftBuildMonitorEvent::Metric { metric, value } => {
format!("metric:{metric}:{value}")
}
});
Ok(())
};
let monitor = NtfsMftBuildMonitor::new_with_observer(None, false, &mut observer);
monitor
.measure(NtfsMftBuildStage::OpenVolume, || Ok(()))
.unwrap();
monitor.add_metric(NtfsMftBuildMetric::ParsedRecords, 2);
monitor.add_metric(NtfsMftBuildMetric::ParsedRecords, 3);
}
assert_eq!(
events,
[
"started:open-volume",
"finished:open-volume",
"metric:parsed-records:2",
"metric:parsed-records:5",
]
);
}
#[test]
fn mft_build_monitor_observer_error_aborts_next_progress_check() {
let mut observer = |_event| Err(RebeccaError::Io(io::Error::other("progress sink failed")));
let monitor = NtfsMftBuildMonitor::new_with_observer(None, false, &mut observer);
monitor.add_metric(NtfsMftBuildMetric::ParsedRecords, 1);
let err = check_mft_build_progress(&ScanCancellationToken::new(), &monitor).unwrap_err();
assert!(err.to_string().contains("progress sink failed"));
}
#[test]
fn sequential_mft_parallel_parse_preserves_chunk_order_and_base_ids() {
let reader = MftRecordReader::new(4, 1);
let chunks = vec![
SequentialMftChunk {
base_record_id: 20,
bytes: vec![0; 8],
mirror: None,
},
SequentialMftChunk {
base_record_id: 10,
bytes: vec![0; 4],
mirror: None,
},
];
let batches =
parse_sequential_mft_chunks(&reader, &ScanCancellationToken::new(), &chunks).unwrap();
let error_record_ids: Vec<_> = batches
.into_iter()
.flat_map(|batch| batch.errors)
.map(|error| error.record_id)
.collect();
assert_eq!(error_record_ids, vec![20, 21, 10]);
}
#[test]
fn sequential_mft_parallel_parse_preserves_cancellation() {
let reader = MftRecordReader::new(4, 1);
let cancellation = ScanCancellationToken::new();
cancellation.cancel();
let err = parse_sequential_mft_chunks(
&reader,
&cancellation,
&[SequentialMftChunk {
base_record_id: 0,
bytes: vec![0; 4],
mirror: None,
}],
)
.unwrap_err();
assert!(matches!(err, RebeccaError::OperationCancelled(_)));
}
#[test]
fn record_source_strategy_returns_first_success() {
let monitor = test_build_monitor();
let source = FakeRecordSource {
label: "primary",
behavior: FakeRecordSourceBehavior::Success("primary-success"),
};
let records = read_mft_records_from_sources(
&[&source],
&NTFS_VOLUME_DATA_BUFFER::default(),
&ScanCancellationToken::new(),
&monitor,
)
.unwrap();
assert_eq!(records.source_label, "primary");
assert_eq!(records.caveats.len(), 1);
assert_eq!(records.caveats[0].code, "primary-success");
}
#[test]
fn record_source_strategy_tries_next_fallback_capable_source() {
let monitor = test_build_monitor();
let unavailable = FakeRecordSource {
label: "sequential",
behavior: FakeRecordSourceBehavior::PlatformUnavailable,
};
let fallback = FakeRecordSource {
label: "fsctl-record",
behavior: FakeRecordSourceBehavior::Success("fallback-success"),
};
let records = read_mft_records_from_sources(
&[&unavailable, &fallback],
&NTFS_VOLUME_DATA_BUFFER::default(),
&ScanCancellationToken::new(),
&monitor,
)
.unwrap();
assert_eq!(records.source_label, "fsctl-record");
assert!(
records
.caveats
.iter()
.any(|caveat| caveat.code == "fallback-success")
);
assert!(records.caveats.iter().any(|caveat| {
caveat.code == "mft-record-source-fallback" && caveat.message.contains("sequential")
}));
}
#[test]
fn record_source_strategy_preserves_cancelled_error() {
let monitor = test_build_monitor();
let cancelled = FakeRecordSource {
label: "sequential",
behavior: FakeRecordSourceBehavior::Cancelled,
};
let fallback = FakeRecordSource {
label: "fsctl-record",
behavior: FakeRecordSourceBehavior::Success("fallback-success"),
};
let err = read_mft_records_from_sources(
&[&cancelled, &fallback],
&NTFS_VOLUME_DATA_BUFFER::default(),
&ScanCancellationToken::new(),
&monitor,
)
.unwrap_err();
assert!(matches!(err, RebeccaError::OperationCancelled(_)));
}
#[test]
fn record_source_strategy_stops_when_build_budget_expires() {
let monitor = NtfsMftBuildMonitor::expired_for_test(Duration::from_secs(1));
let source = FakeRecordSource {
label: "primary",
behavior: FakeRecordSourceBehavior::Success("should-not-run"),
};
let err = read_mft_records_from_sources(
&[&source],
&NTFS_VOLUME_DATA_BUFFER::default(),
&ScanCancellationToken::new(),
&monitor,
)
.unwrap_err();
assert!(matches!(err, RebeccaError::PlatformUnavailable(_)));
assert!(err.to_string().contains("timed out"));
}
#[test]
fn parse_error_caveats_are_sampled_with_summary() {
let mut parse_errors = MftParseErrorCaveats::default();
for record_id in 0..MAX_MFT_PARSE_ERROR_CAVEAT_SAMPLES + 3 {
parse_errors.record(record_id as u64, "invalid signature");
}
let mut caveats = Vec::new();
parse_errors.append_to(&mut caveats);
assert_eq!(caveats.len(), MAX_MFT_PARSE_ERROR_CAVEAT_SAMPLES + 1);
assert_eq!(caveats[0].code, "mft-record-parse-error");
assert!(caveats[0].message.contains("record 0"));
let summary = caveats.last().unwrap();
assert_eq!(summary.code, "mft-record-parse-error-summary");
assert!(summary.message.contains("3 additional"));
}
#[test]
fn estimate_caveats_are_bounded_per_code() {
let measured = MeasuredScan::exact(
ScanReport {
bytes_scanned: 0,
files_scanned: 0,
directories_scanned: 0,
},
ScanBackendKind::WindowsNtfsMftExperimental,
);
let mut caveats: Vec<_> = (0..MAX_MFT_ESTIMATE_CAVEAT_SAMPLES_PER_CODE + 2)
.map(|index| {
ParseCaveat::new(
"multiple-file-names",
format!("record {index} has multiple names"),
)
})
.collect();
caveats.push(ParseCaveat::new(
"attribute-list-present",
"record uses an attribute list",
));
let measured = with_bounded_mft_caveats(measured, caveats);
assert_eq!(
measured
.caveats
.iter()
.filter(|caveat| caveat.code == "multiple-file-names")
.count(),
MAX_MFT_ESTIMATE_CAVEAT_SAMPLES_PER_CODE
);
assert!(measured.caveats.iter().any(|caveat| {
caveat.code == MFT_CAVEAT_SUMMARY_CODE
&& caveat
.message
.contains("2 additional 'multiple-file-names'")
}));
assert!(measured.caveats.iter().any(|caveat| {
caveat.code == "attribute-list-present"
&& caveat.message == "record uses an attribute list"
}));
}
struct FakeRecordSource {
label: &'static str,
behavior: FakeRecordSourceBehavior,
}
enum FakeRecordSourceBehavior {
Success(&'static str),
PlatformUnavailable,
Cancelled,
}
impl MftRecordSource for FakeRecordSource {
fn label(&self) -> &'static str {
self.label
}
fn read_records(
&self,
_volume_data: &NTFS_VOLUME_DATA_BUFFER,
_cancellation: &ScanCancellationToken,
_monitor: &NtfsMftBuildMonitor,
) -> Result<ParsedNtfsRecords> {
match self.behavior {
FakeRecordSourceBehavior::Success(code) => Ok(ParsedNtfsRecords {
source_label: self.label,
records: Vec::new(),
caveats: vec![ParseCaveat::new(code, self.label)],
}),
FakeRecordSourceBehavior::PlatformUnavailable => Err(
RebeccaError::PlatformUnavailable("not available".to_string()),
),
FakeRecordSourceBehavior::Cancelled => {
Err(RebeccaError::OperationCancelled("cancelled".to_string()))
}
}
}
}
#[derive(Default)]
struct FakeIndexStreamSource {
bytes: BTreeMap<u64, u8>,
read_delay: Option<Duration>,
read_count: usize,
max_reads: Option<usize>,
}
impl FakeIndexStreamSource {
fn with_bytes(mut self, offset: u64, bytes: &[u8]) -> Self {
for (index, byte) in bytes.iter().copied().enumerate() {
self.bytes.insert(offset + index as u64, byte);
}
self
}
fn with_read_delay(mut self, delay: Duration) -> Self {
self.read_delay = Some(delay);
self
}
fn with_max_reads(mut self, max_reads: usize) -> Self {
self.max_reads = Some(max_reads);
self
}
}
impl NtfsStreamSource for FakeIndexStreamSource {
type Error = &'static str;
fn read_bytes_at(
&mut self,
volume_offset: u64,
len: usize,
) -> std::result::Result<Vec<u8>, Self::Error> {
if let Some(delay) = self.read_delay {
std::thread::sleep(delay);
}
self.read_count += 1;
let mut bytes = Vec::new();
for index in 0..len {
let Some(byte) = self.bytes.get(&(volume_offset + index as u64)) else {
break;
};
bytes.push(*byte);
}
Ok(bytes)
}
fn should_continue_stream_reads(&self) -> bool {
self.max_reads
.is_none_or(|max_reads| self.read_count < max_reads)
}
}
#[derive(Default)]
struct FakeTargetedRecordResolver {
records: BTreeMap<u64, NtfsParsedRecord>,
reads: Vec<u64>,
}
impl FakeTargetedRecordResolver {
fn with_record(mut self, record: NtfsParsedRecord) -> Self {
self.records.insert(record.reference.record_id, record);
self
}
}
impl TargetedMftRecordResolver for FakeTargetedRecordResolver {
fn resolve_record(
&mut self,
reference: NtfsFileReference,
) -> Result<Option<NtfsParsedRecord>> {
self.reads.push(reference.record_id);
Ok(self.records.get(&reference.record_id).cloned())
}
}
struct CancellingIndexStreamSource {
cancellation: ScanCancellationToken,
}
impl NtfsStreamSource for CancellingIndexStreamSource {
type Error = &'static str;
fn read_bytes_at(
&mut self,
_volume_offset: u64,
_len: usize,
) -> std::result::Result<Vec<u8>, Self::Error> {
self.cancellation.cancel();
Err("cancelled")
}
}
fn test_record_geometry() -> NtfsRecordGeometry {
NtfsRecordGeometry {
record_size: 1024,
sector_size: 512,
bytes_per_cluster: 4096,
max_record_count: 16,
}
}
fn test_build_monitor() -> NtfsMftBuildMonitor<'static> {
NtfsMftBuildMonitor::new(None, false)
}
fn disk_map_options(
top_limit: usize,
max_depth: Option<usize>,
group_kinds: Vec<DiskMapGroupKind>,
) -> DiskMapBackendOptions {
DiskMapBackendOptions {
top_limit,
top_sort: DiskMapSortField::Logical,
entry_filter: Default::default(),
max_depth,
group_kinds,
group_limit: 20,
group_now: UNIX_EPOCH,
group_sort: DiskMapSortField::Logical,
}
}
fn assert_group_metrics(
groups: &[DiskMapGroup],
kind: DiskMapGroupKind,
key: &str,
logical_bytes: u64,
files: u64,
) {
let group = groups
.iter()
.find(|group| group.kind == kind && group.key == key)
.unwrap_or_else(|| panic!("missing group {}:{key}", kind.label()));
assert_eq!(group.metrics.logical_bytes, logical_bytes);
assert_eq!(group.metrics.files, files);
}
fn parsed_directory_with_index_allocation(record_id: u64, name: &str) -> NtfsParsedRecord {
NtfsParsedRecord {
reference: NtfsFileReference::known(record_id, record_id as u16),
base_reference: None,
in_use: true,
is_directory: true,
is_reparse_point: false,
attributes: Vec::new(),
attribute_list_entries: Vec::new(),
names: vec![parsed_file_name(record_id, name, FILE_ATTRIBUTE_DIRECTORY)],
attribute_streams: vec![NtfsAttributeStream {
attribute_type: AttributeType::IndexAllocation,
attribute_id: 0,
name: Some("$I30".to_string()),
non_resident: true,
flags: 0,
lowest_vcn: Some(0),
highest_vcn: Some(0),
logical_size: RECORD_SIZE as u64,
allocated_size: Some(RECORD_SIZE as u64),
initialized_size: Some(RECORD_SIZE as u64),
data_runs: vec![NtfsDataRun {
starting_vcn: 0,
cluster_count: 1,
lcn: Some(0x80),
}],
}],
directory_indexes: vec![NtfsDirectoryIndex {
name: "$I30".to_string(),
attribute_id: 0,
indexed_attribute: AttributeType::FileName,
index_record_size: RECORD_SIZE as u32,
root_entries: vec![NtfsIndexEntry {
directory_entry: None,
child_vcn: Some(0),
is_last: true,
}],
}],
directory_entries: Vec::new(),
caveats: Vec::new(),
}
}
fn parsed_file(record_id: u64, parent_id: u64, name: &str, bytes: u64) -> NtfsParsedRecord {
NtfsParsedRecord {
reference: NtfsFileReference::known(record_id, record_id as u16),
base_reference: None,
in_use: true,
is_directory: false,
is_reparse_point: false,
attributes: Vec::new(),
attribute_list_entries: Vec::new(),
names: vec![parsed_file_name(parent_id, name, 0)],
attribute_streams: vec![NtfsAttributeStream {
attribute_type: AttributeType::Data,
attribute_id: 0,
name: None,
non_resident: false,
flags: 0,
lowest_vcn: None,
highest_vcn: None,
logical_size: bytes,
allocated_size: Some(bytes),
initialized_size: Some(bytes),
data_runs: Vec::new(),
}],
directory_indexes: Vec::new(),
directory_entries: Vec::new(),
caveats: Vec::new(),
}
}
fn parsed_file_name(parent_id: u64, name: &str, file_attributes: u32) -> NtfsFileName {
NtfsFileName {
parent: NtfsFileReference::known(parent_id, parent_id as u16),
namespace: FileNameNamespace::Win32,
name: name.to_string(),
attribute_id: Some(0),
attribute_name: None,
lowest_vcn: None,
modified_windows_filetime: 0,
allocated_size: 0,
real_size: 0,
file_attributes,
}
}
fn ntfs_volume_capabilities(device_path: &str, volume_serial: u64) -> NtfsVolumeCapabilities {
let drive = device_path
.strip_prefix("\\\\.\\")
.unwrap_or(device_path)
.trim_end_matches(':');
NtfsVolumeCapabilities {
root_path: std::path::PathBuf::from(format!("{drive}:\\")),
device_path: device_path.to_string(),
mft_data_path: format!("\\\\?\\{drive}:\\$MFT::$DATA"),
volume_serial,
}
}
fn ntfs_volume_fingerprint(
device_path: &str,
volume_serial: u64,
record_size: u32,
sector_size: u32,
bytes_per_cluster: u32,
mft_valid_data_length: i64,
) -> NtfsVolumeIndexFingerprint {
let capabilities = ntfs_volume_capabilities(device_path, volume_serial);
let mut volume_data = ntfs_volume_data(
record_size,
sector_size,
bytes_per_cluster,
mft_valid_data_length,
);
volume_data.MftStartLcn = 12;
volume_data.Mft2StartLcn = 34;
let geometry =
NtfsRecordGeometry::from_volume_data(&capabilities.device_path, &volume_data).unwrap();
NtfsVolumeIndexFingerprint::from_volume_data(&capabilities, &volume_data, geometry)
}
fn usn_state(journal_id: u64, first_usn: u64, next_usn: u64) -> ScanCacheUsnJournalState {
ScanCacheUsnJournalState {
journal_id,
first_usn,
next_usn,
}
}
fn fixture_mft_index() -> MftIndex {
MftIndex::from_parsed_records(vec![
parsed_directory_with_index_allocation(5, "root"),
parsed_file(6, 5, "large.bin", 10),
parsed_file(7, 5, "small.txt", 3),
])
}
fn nested_fixture_mft_index() -> MftIndex {
MftIndex::from_parsed_records(vec![
parsed_directory(5, 5, "root"),
parsed_directory(8, 5, "target"),
parsed_file(9, 8, "inside.txt", 1),
parsed_directory(20, 5, "outside"),
parsed_file(30, 20, "outside.txt", 1),
])
}
fn parsed_directory(record_id: u64, parent_id: u64, name: &str) -> NtfsParsedRecord {
let mut record = parsed_directory_with_index_allocation(record_id, name);
record.names = vec![parsed_file_name(parent_id, name, FILE_ATTRIBUTE_DIRECTORY)];
record.attribute_streams.clear();
record.directory_indexes.clear();
record
}
fn usn_read_buffer(next_usn: i64, records: Vec<Vec<u8>>) -> Vec<u8> {
let mut raw = next_usn.to_le_bytes().to_vec();
for record in records {
raw.extend(record);
}
raw
}
fn usn_record_v2(
file_reference: NtfsFileReference,
parent_reference: NtfsFileReference,
usn: i64,
) -> Vec<u8> {
let record_len = offset_of!(USN_RECORD_V2, FileName);
let mut raw = vec![0_u8; record_len];
raw[0..4].copy_from_slice(&(record_len as u32).to_le_bytes());
raw[offset_of!(USN_RECORD_V2, MajorVersion)..][..2].copy_from_slice(&2_u16.to_le_bytes());
raw[offset_of!(USN_RECORD_V2, MinorVersion)..][..2].copy_from_slice(&0_u16.to_le_bytes());
raw[offset_of!(USN_RECORD_V2, FileReferenceNumber)..][..8]
.copy_from_slice(&file_reference_number(file_reference).to_le_bytes());
raw[offset_of!(USN_RECORD_V2, ParentFileReferenceNumber)..][..8]
.copy_from_slice(&file_reference_number(parent_reference).to_le_bytes());
raw[offset_of!(USN_RECORD_V2, Usn)..][..8].copy_from_slice(&usn.to_le_bytes());
raw[offset_of!(USN_RECORD_V2, FileNameLength)..][..2].copy_from_slice(&0_u16.to_le_bytes());
raw[offset_of!(USN_RECORD_V2, FileNameOffset)..][..2]
.copy_from_slice(&(record_len as u16).to_le_bytes());
raw
}
fn ntfs_volume_data(
record_size: u32,
sector_size: u32,
bytes_per_cluster: u32,
mft_valid_data_length: i64,
) -> NTFS_VOLUME_DATA_BUFFER {
NTFS_VOLUME_DATA_BUFFER {
BytesPerFileRecordSegment: record_size,
BytesPerSector: sector_size,
BytesPerCluster: bytes_per_cluster,
MftValidDataLength: mft_valid_data_length,
..Default::default()
}
}
const RECORD_SIZE: usize = 1024;
const SECTOR_SIZE: usize = 512;
const FILE_ATTRIBUTE_NORMAL: u32 = 0x0000_0080;
const FILE_ATTRIBUTE_DIRECTORY: u32 = 0x0000_0010;
fn index_allocation_record(vcn: u64, entries: Vec<Vec<u8>>) -> Vec<u8> {
let mut raw_entries = Vec::new();
for entry in entries {
raw_entries.extend_from_slice(&entry);
}
let mut record = vec![0_u8; RECORD_SIZE];
let usa_offset = 0x28;
let index_header_offset = 0x18;
let entries_offset = 0x20;
let entries_start = index_header_offset + entries_offset;
let index_size = entries_offset + raw_entries.len();
record[0..4].copy_from_slice(b"INDX");
put_u16(&mut record, 4, usa_offset as u16);
put_u16(&mut record, 6, 3);
put_u64(&mut record, 16, vcn);
put_u32(&mut record, index_header_offset, entries_offset as u32);
put_u32(&mut record, index_header_offset + 4, index_size as u32);
put_u32(&mut record, index_header_offset + 8, index_size as u32);
record[entries_start..entries_start + raw_entries.len()].copy_from_slice(&raw_entries);
apply_test_fixup_at(&mut record, usa_offset);
record
}
fn index_allocation_entry(
child_reference: u64,
parent_reference: u64,
name: &str,
file_attributes: u32,
) -> Vec<u8> {
index_allocation_entry_with_namespace(
child_reference,
parent_reference,
name,
file_attributes,
1,
)
}
fn index_allocation_entry_with_namespace(
child_reference: u64,
parent_reference: u64,
name: &str,
file_attributes: u32,
namespace: u8,
) -> Vec<u8> {
let file_name =
file_name_value_with_namespace(parent_reference, name, file_attributes, namespace);
let entry_len = align8(16 + file_name.len());
let mut entry = vec![0_u8; entry_len];
put_u64(&mut entry, 0, child_reference);
put_u16(&mut entry, 8, entry_len as u16);
put_u16(&mut entry, 10, file_name.len() as u16);
entry[16..16 + file_name.len()].copy_from_slice(&file_name);
entry
}
fn index_allocation_last_entry() -> Vec<u8> {
let mut entry = vec![0_u8; 16];
put_u16(&mut entry, 8, 16);
put_u16(&mut entry, 12, 0x0002);
entry
}
fn file_name_value_with_namespace(
parent_reference: u64,
name: &str,
file_attributes: u32,
namespace: u8,
) -> Vec<u8> {
let name_utf16 = name.encode_utf16().collect::<Vec<_>>();
let mut value = vec![0_u8; 66 + (name_utf16.len() * 2)];
put_u64(&mut value, 0, parent_reference);
put_u32(&mut value, 56, file_attributes);
value[64] = name_utf16.len() as u8;
value[65] = namespace;
for (index, character) in name_utf16.iter().enumerate() {
put_u16(&mut value, 66 + (index * 2), *character);
}
value
}
fn file_reference(record_id: u64, sequence_number: u16) -> u64 {
((sequence_number as u64) << 48) | (record_id & 0x0000_FFFF_FFFF_FFFF)
}
fn apply_test_fixup_at(record: &mut [u8], usa_offset: usize) {
let update_sequence = 0xBBAA_u16;
let sector_count = record.len() / SECTOR_SIZE;
put_u16(record, usa_offset, update_sequence);
for sector_index in 0..sector_count {
let tail = ((sector_index + 1) * SECTOR_SIZE) - 2;
let original = u16::from_le_bytes([record[tail], record[tail + 1]]);
put_u16(record, usa_offset + ((sector_index + 1) * 2), original);
put_u16(record, tail, update_sequence);
}
}
fn put_u16(bytes: &mut [u8], offset: usize, value: u16) {
bytes[offset..offset + 2].copy_from_slice(&value.to_le_bytes());
}
fn put_u32(bytes: &mut [u8], offset: usize, value: u32) {
bytes[offset..offset + 4].copy_from_slice(&value.to_le_bytes());
}
fn put_u64(bytes: &mut [u8], offset: usize, value: u64) {
bytes[offset..offset + 8].copy_from_slice(&value.to_le_bytes());
}
fn align8(value: usize) -> usize {
(value + 7) & !7
}
fn retrieval_pointer_buffer(starting_vcn: i64, extents: &[(i64, i64)]) -> Vec<u8> {
let header_size = std::mem::offset_of!(RETRIEVAL_POINTERS_BUFFER, Extents);
let extent_size = std::mem::size_of::<RETRIEVAL_POINTERS_BUFFER_0>();
let mut buffer = vec![0_u8; header_size + (extent_size * extents.len())];
unsafe {
std::ptr::write_unaligned(
buffer.as_mut_ptr().cast::<RETRIEVAL_POINTERS_BUFFER>(),
RETRIEVAL_POINTERS_BUFFER {
ExtentCount: extents.len() as u32,
StartingVcn: starting_vcn,
Extents: [RETRIEVAL_POINTERS_BUFFER_0::default()],
},
);
for (index, (next_vcn, lcn)) in extents.iter().copied().enumerate() {
std::ptr::write_unaligned(
buffer
.as_mut_ptr()
.add(header_size + (index * extent_size))
.cast::<RETRIEVAL_POINTERS_BUFFER_0>(),
RETRIEVAL_POINTERS_BUFFER_0 {
NextVcn: next_vcn,
Lcn: lcn,
},
);
}
}
buffer
}
}