use super::{
FileAggregate, FileReport, MAX_JSONL_LINES, MeasurementLimits, SamplingReport, codecs,
u64_from_usize,
};
use cap_fs_ext::{DirEntryExt, DirExt, FollowSymlinks, OpenOptionsFollowExt};
use cap_std::{
ambient_authority,
fs::{Dir as CapDir, File as CapFile, Metadata as CapMetadata, OpenOptions as CapOpenOptions},
};
use std::{
ffi::OsStr,
io::{self, BufRead, BufReader, Read},
path::{Path, PathBuf},
time::SystemTime,
};
pub(super) const MAX_SESSION_FILE_BYTES: u64 = 64 * 1024 * 1024;
pub(super) const MAX_JSONL_LINE_BYTES: usize = 1024 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct FileSignature {
len: u64,
modified: SystemTime,
#[cfg(any(unix, windows))]
dev: u64,
#[cfg(any(unix, windows))]
ino: u64,
#[cfg(any(unix, windows))]
nlink: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct DiscoverySignature {
len: u64,
modified: SystemTime,
#[cfg(any(unix, windows))]
dev: u64,
#[cfg(any(unix, windows))]
ino: u64,
#[cfg(any(unix, windows))]
nlink: u64,
}
#[derive(Debug)]
pub(super) struct FileTarget {
path: PathBuf,
file: CapFile,
discovery_signature: DiscoverySignature,
opened_signature: FileSignature,
}
#[derive(Debug, Default)]
pub(super) struct Discovery {
pub(super) targets: Vec<FileTarget>,
pub(super) files: FileReport,
pub(super) sampling: SamplingReport,
pub(super) root_missing: bool,
pub(super) unsafe_storage: bool,
}
#[derive(Debug)]
pub(super) enum FileScanResult {
Scanned {
aggregate: Box<FileAggregate>,
bytes: u64,
},
Unreadable {
bytes: u64,
},
Unstable {
bytes: u64,
},
}
pub(super) fn discover_files(root: &Path, limits: MeasurementLimits) -> Discovery {
let mut discovery = Discovery {
sampling: SamplingReport {
max_files: u64_from_usize(limits.max_files),
max_bytes: limits.max_bytes,
..SamplingReport::default()
},
..Discovery::default()
};
let root_directory = match open_root_directory(root) {
Ok(directory) => directory,
Err(error) if error.kind() == io::ErrorKind::NotFound => {
discovery.root_missing = true;
return discovery;
}
Err(_) => {
discovery.unsafe_storage = true;
return discovery;
}
};
let mut selector = CandidateSelector::new(limits);
collect_directory_entries(
&root_directory,
Path::new("primary"),
false,
&mut selector,
&mut discovery,
);
match root_directory.open_dir_nofollow("subagents") {
Ok(subagents) => match subagents.dir_metadata() {
Ok(metadata) if !metadata_is_unsafe(&metadata) && metadata.is_dir() => {
collect_directory_entries(
&subagents,
Path::new("subagents"),
true,
&mut selector,
&mut discovery,
);
}
Ok(_) | Err(_) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
}
},
Err(error) if error.kind() == io::ErrorKind::NotFound => {}
Err(_) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
}
}
let (targets, sampling) = selector.finish();
discovery.files.considered = sampling.discovered_files;
discovery.files.sampled_out = sampling.sampled_out_files;
discovery.sampling = sampling;
discovery.targets = targets;
discovery
}
#[derive(Debug)]
struct CandidateSelector {
max_files: usize,
max_bytes: u64,
targets: Vec<FileTarget>,
discovered_files: u64,
discovered_bytes: u64,
}
impl CandidateSelector {
fn new(limits: MeasurementLimits) -> Self {
Self {
max_files: limits.max_files,
max_bytes: limits.max_bytes,
targets: Vec::with_capacity(limits.max_files),
discovered_files: 0,
discovered_bytes: 0,
}
}
fn consider(&mut self, target: FileTarget) {
let size = target.opened_signature.len;
self.discovered_files = self.discovered_files.saturating_add(1);
self.discovered_bytes = self.discovered_bytes.saturating_add(size);
if size > self.max_bytes {
return;
}
if self.targets.len() < self.max_files {
self.targets.push(target);
return;
}
let Some(worst_index) = self
.targets
.iter()
.enumerate()
.max_by(|(_, left), (_, right)| compare_targets(left, right))
.map(|(index, _)| index)
else {
return;
};
if compare_targets(&target, &self.targets[worst_index]).is_lt() {
self.targets[worst_index] = target;
}
}
fn finish(mut self) -> (Vec<FileTarget>, SamplingReport) {
self.targets.sort_by(compare_targets);
let mut selected = Vec::with_capacity(self.targets.len());
let mut planned_bytes = 0u64;
for target in self.targets {
let size = target.opened_signature.len;
let remaining = self.max_bytes.saturating_sub(planned_bytes);
if size <= remaining {
planned_bytes = planned_bytes.saturating_add(size);
selected.push(target);
}
}
let selected_files = u64_from_usize(selected.len());
let sampled_out_files = self.discovered_files.saturating_sub(selected_files);
let sampled_out_bytes = self.discovered_bytes.saturating_sub(planned_bytes);
(
selected,
SamplingReport {
max_files: u64_from_usize(self.max_files),
max_bytes: self.max_bytes,
discovered_files: self.discovered_files,
selected_files,
sampled_out_files,
discovered_bytes: self.discovered_bytes,
planned_bytes,
scanned_bytes: 0,
sampled_out_bytes,
unscanned_selected_bytes: 0,
},
)
}
}
fn compare_targets(left: &FileTarget, right: &FileTarget) -> std::cmp::Ordering {
right
.opened_signature
.modified
.cmp(&left.opened_signature.modified)
.then_with(|| left.path.cmp(&right.path))
}
fn collect_directory_entries(
directory: &CapDir,
scope: &Path,
is_subagents: bool,
selector: &mut CandidateSelector,
discovery: &mut Discovery,
) {
let iterator = match directory.entries() {
Ok(iterator) => iterator,
Err(_) => {
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
if is_subagents {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
}
return;
}
};
for entry in iterator {
let entry = match entry {
Ok(entry) => entry,
Err(_) => {
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
};
let name = entry.file_name();
if !is_subagents && name == OsStr::new("subagents") {
continue;
}
if !is_active_jsonl_name(&name) {
continue;
}
let file_type = match entry.file_type() {
Ok(file_type) => file_type,
Err(_) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
};
if file_type.is_symlink() || !file_type.is_file() {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
let metadata = match entry.full_metadata() {
Ok(metadata) => metadata,
Err(_) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
};
let discovery_signature = match discovery_signature_from_metadata(&metadata) {
Ok(Some(signature)) => signature,
Ok(None) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
Err(_) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
};
let file = match open_child_read_only(directory, &name) {
Ok(file) => file,
Err(_) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
};
let opened_signature = match file_signature_from_handle(&file) {
Ok(Some(signature)) => signature,
Ok(None) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
Err(_) => {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.unreadable = discovery.files.unreadable.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
};
if !discovery_signature_matches_opened(discovery_signature, opened_signature) {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.unsafe_storage = true;
continue;
}
if opened_signature.len > MAX_SESSION_FILE_BYTES {
discovery.files.skipped = discovery.files.skipped.saturating_add(1);
discovery.files.limit_hit = discovery.files.limit_hit.saturating_add(1);
continue;
}
selector.consider(FileTarget {
path: scope.join(Path::new(&name)),
file,
discovery_signature,
opened_signature,
});
}
}
fn is_active_jsonl_name(name: &OsStr) -> bool {
let Some(name) = name.to_str() else {
return false;
};
let Some(id) = name.strip_suffix(".jsonl") else {
return false;
};
!id.is_empty()
&& id
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || byte == b'_' || byte == b'-')
}
#[cfg(any(unix, windows))]
fn open_root_directory(root: &Path) -> io::Result<CapDir> {
use std::path::Component;
let mut components = root.components();
#[cfg(unix)]
let mut directory = {
if !matches!(components.next(), Some(Component::RootDir)) {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"session root must be an absolute directory path",
));
}
CapDir::open_ambient_dir(Path::new("/"), ambient_authority())?
};
#[cfg(windows)]
let mut directory = {
let Some(Component::Prefix(prefix)) = components.next() else {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"session root must be an absolute directory path",
));
};
if !matches!(components.next(), Some(Component::RootDir)) {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"session root must be an absolute directory path",
));
}
let mut anchor = PathBuf::from(prefix.as_os_str());
anchor.push(Path::new("\\"));
CapDir::open_ambient_dir(anchor, ambient_authority())?
};
validate_directory_handle(&directory)?;
let mut component_count = 0usize;
for component in components {
let Component::Normal(name) = component else {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"session root contains an unsafe path component",
));
};
directory = open_directory_child_nofollow(&directory, name)?;
component_count = component_count.saturating_add(1);
}
if component_count == 0 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"session root must have a final directory name",
));
}
Ok(directory)
}
#[cfg(any(unix, windows))]
fn open_directory_child_nofollow(parent: &CapDir, name: &OsStr) -> io::Result<CapDir> {
let directory = parent.open_dir_nofollow(Path::new(name))?;
validate_directory_handle(&directory)?;
Ok(directory)
}
#[cfg(any(unix, windows))]
fn validate_directory_handle(directory: &CapDir) -> io::Result<()> {
let metadata = directory.dir_metadata()?;
if metadata_is_unsafe(&metadata) || !metadata.is_dir() {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"session root is not a safe directory",
));
}
Ok(())
}
#[cfg(not(any(unix, windows)))]
fn open_root_directory(_root: &Path) -> io::Result<CapDir> {
Err(io::Error::new(
io::ErrorKind::Unsupported,
"offline measurement requires no-follow directory capabilities on this platform",
))
}
fn open_child_read_only(directory: &CapDir, name: &OsStr) -> io::Result<CapFile> {
let mut options = CapOpenOptions::new();
options.read(true).follow(FollowSymlinks::No);
directory.open_with(Path::new(name), &options)
}
pub(super) fn scan_file(target: &FileTarget, model: &str) -> FileScanResult {
let expected_bytes = target.opened_signature.len;
let before = match file_signature_from_handle(&target.file) {
Ok(Some(signature)) => signature,
Ok(None) => {
return FileScanResult::Unstable {
bytes: expected_bytes,
};
}
Err(_) => {
return FileScanResult::Unreadable {
bytes: expected_bytes,
};
}
};
if !discovery_signature_matches_opened(target.discovery_signature, before)
|| !signature_is_unchanged(target.opened_signature, Some(before))
|| before.len > MAX_SESSION_FILE_BYTES
{
return FileScanResult::Unstable {
bytes: expected_bytes,
};
}
let mut aggregate = FileAggregate::default();
if scan_file_lines(&target.file, before.len, model, &mut aggregate).is_err() {
return FileScanResult::Unreadable {
bytes: expected_bytes,
};
}
let after = match file_signature_from_handle(&target.file) {
Ok(Some(signature)) => signature,
Ok(None) => {
return FileScanResult::Unstable {
bytes: expected_bytes,
};
}
Err(_) => {
return FileScanResult::Unreadable {
bytes: expected_bytes,
};
}
};
if !discovery_signature_matches_opened(target.discovery_signature, after)
|| !signature_is_unchanged(before, Some(after))
{
return FileScanResult::Unstable {
bytes: expected_bytes,
};
}
FileScanResult::Scanned {
aggregate: Box::new(aggregate),
bytes: before.len,
}
}
fn signature_is_unchanged(before: FileSignature, after: Option<FileSignature>) -> bool {
after == Some(before)
}
#[derive(Debug)]
struct BoundedLine {
bytes: Vec<u8>,
oversized: bool,
}
fn read_bounded_line<R: BufRead>(
reader: &mut R,
max_bytes: usize,
) -> io::Result<Option<BoundedLine>> {
let mut bytes = Vec::with_capacity(max_bytes.min(8192));
let mut length = 0usize;
let mut saw_bytes = false;
let mut oversized = false;
loop {
let buffer = reader.fill_buf()?;
if buffer.is_empty() {
if !saw_bytes {
return Ok(None);
}
return Ok(Some(BoundedLine { bytes, oversized }));
}
let newline = buffer.iter().position(|byte| *byte == b'\n');
let consumed = newline.map_or(buffer.len(), |index| index.saturating_add(1));
let content_len = newline.unwrap_or(buffer.len());
if !oversized {
let Some(next_length) = length.checked_add(content_len) else {
oversized = true;
bytes.clear();
length = max_bytes.saturating_add(1);
reader.consume(consumed);
saw_bytes = true;
if newline.is_some() {
return Ok(Some(BoundedLine { bytes, oversized }));
}
continue;
};
if next_length > max_bytes {
oversized = true;
bytes.clear();
length = max_bytes.saturating_add(1);
} else {
bytes.extend_from_slice(&buffer[..consumed]);
length = next_length;
}
}
reader.consume(consumed);
saw_bytes = true;
if newline.is_some() {
return Ok(Some(BoundedLine { bytes, oversized }));
}
}
}
fn file_signature_from_handle(file: &CapFile) -> io::Result<Option<FileSignature>> {
let metadata = file.metadata()?;
if metadata_is_unsafe(&metadata)
|| !metadata.file_type().is_file()
|| has_unexpected_hard_links(&metadata)
{
return Ok(None);
}
#[cfg(any(unix, windows))]
{
use cap_fs_ext::MetadataExt;
Ok(Some(FileSignature {
len: metadata.len(),
modified: metadata.modified()?.into_std(),
dev: metadata.dev(),
ino: metadata.ino(),
nlink: metadata.nlink(),
}))
}
#[cfg(not(any(unix, windows)))]
{
Err(io::Error::new(
io::ErrorKind::Unsupported,
"offline measurement requires stable file identity on this platform",
))
}
}
fn metadata_is_unsafe(metadata: &CapMetadata) -> bool {
if metadata.file_type().is_symlink() {
return true;
}
#[cfg(windows)]
{
use cap_fs_ext::OsMetadataExt;
const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x400;
metadata.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT != 0
}
#[cfg(not(windows))]
{
false
}
}
fn has_unexpected_hard_links(metadata: &CapMetadata) -> bool {
#[cfg(any(unix, windows))]
{
use cap_fs_ext::MetadataExt;
metadata.nlink() != 1
}
#[cfg(not(any(unix, windows)))]
{
let _ = metadata;
true
}
}
fn discovery_signature_from_metadata(
metadata: &CapMetadata,
) -> io::Result<Option<DiscoverySignature>> {
if metadata_is_unsafe(metadata) || !metadata.is_file() {
return Ok(None);
}
let modified = metadata.modified()?.into_std();
#[cfg(any(unix, windows))]
{
use cap_fs_ext::MetadataExt;
let nlink = metadata.nlink();
if nlink != 1 {
return Ok(None);
}
Ok(Some(DiscoverySignature {
len: metadata.len(),
modified,
dev: metadata.dev(),
ino: metadata.ino(),
nlink,
}))
}
#[cfg(not(any(unix, windows)))]
{
Err(io::Error::new(
io::ErrorKind::Unsupported,
"offline measurement requires stable file identity on this platform",
))
}
}
fn discovery_signature_matches_opened(
discovery: DiscoverySignature,
opened: FileSignature,
) -> bool {
if discovery.len != opened.len || discovery.modified != opened.modified {
return false;
}
#[cfg(any(unix, windows))]
{
discovery.dev == opened.dev
&& discovery.ino == opened.ino
&& discovery.nlink == opened.nlink
}
#[cfg(not(any(unix, windows)))]
{
false
}
}
fn scan_file_lines(
file: &CapFile,
snapshot_len: u64,
model: &str,
aggregate: &mut FileAggregate,
) -> io::Result<()> {
let bounded_file = file
.try_clone()?
.take(snapshot_len.min(MAX_SESSION_FILE_BYTES));
let mut reader = BufReader::new(bounded_file);
for _ in 0..MAX_JSONL_LINES {
let Some(line) = read_bounded_line(&mut reader, MAX_JSONL_LINE_BYTES)? else {
return Ok(());
};
aggregate.jsonl.lines = aggregate.jsonl.lines.saturating_add(1);
if line.oversized {
aggregate.jsonl.oversized = aggregate.jsonl.oversized.saturating_add(1);
aggregate.limit_hit = true;
continue;
}
codecs::process_jsonl_line(&line.bytes, model, aggregate);
}
if !reader.fill_buf()?.is_empty() {
aggregate.jsonl.limit_hit = aggregate.jsonl.limit_hit.saturating_add(1);
aggregate.limit_hit = true;
}
Ok(())
}