Fan the per-entry listing syscalls out across bounded threads

Sniffing every entry's leading bytes made a 5000-entry listing cost 37.5 ms
through the router. `list_dir` now splits into two passes: walking the
directory stream, which is inherently sequential, and the per-entry
syscalls, which are not.

Pass 2 does `metadata` and the content sniff together on worker threads.
Measured warm on 22 cores, best of 5 through the real router:

    entries   before   after
        200    1.2 ms   1.3 ms   (below the threshold, serial)
       1000    5.8 ms   3.7 ms
       5000   37.5 ms   9.6 ms

Moving `metadata` into the parallel pass mattered as much as parallelising
the sniff: it was the largest remaining serial step, and the placeholder
`Entry` fields double as its failure fallbacks, so the broken-symlink
behaviour is unchanged.

Fan-out is bounded by a process-wide `SNIFF_BUDGET` sized to
`available_parallelism()`, claimed with `fetch_update` and returned on
drop. It degrades instead of queueing: one large listing may take the whole
budget, and a concurrent one then gets fewer threads, or none and runs
serially. So the worst case across all in-flight listings is
`available_parallelism()` extra threads, not one set per request. Below 256
entries nothing is spawned, because there the thread cost is the work.

Sorting happens after pass 2, so chunking cannot affect the order. The test
gives every file a distinct length and checks `is_dir`, `size`, `mtime` and
`kind` per row against a fresh single-threaded read, since a chunk-boundary
mistake would otherwise hand one row another row's metadata. A second test
pins the budget arithmetic; both take a mutex because the budget is
process-wide and cargo runs them concurrently.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
AuthorKonata <konata@posteo.jp>
Date
Commit7abefd3eb8ce06775f3cb040e1adf1579218ff94
Parent8100ce9
1 file changed, 198 insertions(+), 28 deletions(-)
▾Mserver/src/fs.rs
@@ -102,26 +102,28 @@ pub fn list_dir(dir: &Path) -> Result<Vec<Entry>, FsError> {
_ => FsError::Forbidden,
})?;
let mut entries = Vec::new();
// Pass 1: names only. Walking the directory stream is inherently
// sequential, but it is also the only part that has to be: every field
// below comes from a per-entry syscall, which pass 2 can do in parallel.
// The values here are the fallbacks used when that syscall fails.
let mut rows: Vec<(Entry, PathBuf)> = Vec::new();
for e in rd.flatten() {
let name = e.file_name().to_string_lossy().into_owned();
let path = e.path();
// Follows symlinks; a broken link shows up as an empty file.
let meta = std::fs::metadata(&path);
let (is_dir, size, mtime) = match meta {
Ok(m) => (m.is_dir(), m.len(), mtime_str(&m)),
Err(_) => (false, 0, "1970-01-01T00:00:00Z".to_string()),
let entry = Entry {
name: e.file_name().to_string_lossy().into_owned(),
is_dir: false,
size: 0,
mtime: EPOCH_MTIME.to_string(),
kind: FileKind::Binary,
};
entries.push(Entry {
name,
is_dir,
size,
mtime,
kind: detect_kind(&path, is_dir),
});
rows.push((entry, e.path()));
}
// Folders first, then case-insensitive name.
// Pass 2: the per-entry syscalls — metadata plus the content sniff.
describe_rows(&mut rows);
let mut entries: Vec<Entry> = rows.into_iter().map(|(e, _)| e).collect();
// Folders first, then case-insensitive name. Sorting after pass 2 means
// the parallel fan-out cannot affect the order.
entries.sort_by(|a, b| {
b.is_dir
.cmp(&a.is_dir)
@@ -137,22 +139,109 @@ pub fn list_dir(dir: &Path) -> Result<Vec<Entry>, FsError> {
/// How many leading bytes we read to classify a file. Every magic number
/// `infer` knows lives in the first few dozen bytes; 256 also gives the
/// text/binary heuristic enough to work with. Measured: ~4.6 µs per file,
/// against ~1.4 µs for the `metadata` call we already make.
/// text/binary heuristic enough to work with. Measured at ~4.6 µs per file,
/// against ~1.4 µs for the `metadata` call in the same pass.
const SNIFF_BYTES: usize = 256;
/// Entries per thread, and the point below which parallelism is not worth it.
/// Measured on a 22-core machine: at 50 entries fan-out is a wash (thread
/// spawn costs about as much as the work), at 200 it is already 2x.
const SNIFF_CHUNK: usize = 256;
/// Process-wide ceiling on threads spawned for sniffing, so many concurrent
/// listings of large directories cannot multiply into a thread explosion.
/// One listing alone can use the whole budget; the next one degrades to fewer
/// threads, and eventually to serial, instead of queueing.
static SNIFF_BUDGET: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
fn sniff_budget_total() -> usize {
static TOTAL: std::sync::LazyLock<usize> = std::sync::LazyLock::new(|| {
std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1)
});
*TOTAL
}
/// Claimed helper threads, returned to [`SNIFF_BUDGET`] on drop.
struct SniffPermit(usize);
impl SniffPermit {
/// Claim up to `want` helper threads, or fewer when the budget is thin.
fn claim(want: usize) -> Self {
use std::sync::atomic::Ordering;
let total = sniff_budget_total();
let mut granted = 0;
let _ = SNIFF_BUDGET.fetch_update(Ordering::AcqRel, Ordering::Acquire, |in_flight| {
granted = want.min(total.saturating_sub(in_flight));
(granted > 0).then_some(in_flight + granted)
});
Self(granted)
}
}
impl Drop for SniffPermit {
fn drop(&mut self) {
if self.0 > 0 {
SNIFF_BUDGET.fetch_sub(self.0, std::sync::atomic::Ordering::AcqRel);
}
}
}
/// Fill in each row from its own syscalls — `metadata` and the content sniff —
/// fanning them out across threads for big directories.
///
/// The calling thread takes a chunk too, so `n` helper threads process `n + 1`
/// chunks and a zero-thread grant is simply the serial path.
fn describe_rows(rows: &mut [(Entry, PathBuf)]) {
fn describe(rows: &mut [(Entry, PathBuf)]) {
for (entry, path) in rows {
// Follows symlinks; a broken link keeps the caller's fallbacks and
// so shows up as an empty file.
if let Ok(m) = std::fs::metadata(&path) {
entry.is_dir = m.is_dir();
entry.size = m.len();
entry.mtime = mtime_str(&m);
}
entry.kind = detect_kind(path, entry.is_dir);
}
}
if rows.len() < SNIFF_CHUNK {
describe(rows);
return;
}
// One chunk per helper thread plus one for this thread.
let want = rows.len().div_ceil(SNIFF_CHUNK).saturating_sub(1);
let permit = SniffPermit::claim(want);
if permit.0 == 0 {
describe(rows);
return;
}
let chunk = rows.len().div_ceil(permit.0 + 1);
std::thread::scope(|s| {
let mut rest = rows;
// Hand every chunk but the last to a helper thread.
while rest.len() > chunk {
let (head, tail) = rest.split_at_mut(chunk);
s.spawn(|| describe(head));
rest = tail;
}
describe(rest);
});
}
/// Classify a directory entry by reading its first [`SNIFF_BYTES`] bytes.
///
/// Blocking — only called from `list_dir` (itself under `spawn_blocking`).
/// An unreadable file is reported as [`FileKind::Binary`] rather than failing
/// the whole listing.
/// Blocking — called from `list_dir` (itself under `spawn_blocking`), on
/// several threads at once for large directories. An unreadable file is
/// reported as [`FileKind::Binary`] rather than failing the whole listing.
///
// ponytail: one open() per entry, serially. Measured ~4.6 µs/file warm
// (23 ms for 5000), against ~1.4 µs for the metadata call — fine on local
// disk. The ceiling is cold cache and network filesystems (NFS/SMB), where
// this becomes a round-trip per entry. If that shows up: sniff only entries
// whose size is non-zero and cache by (dev, ino, mtime), or fan the sniffs
// out across the blocking pool.
// ponytail: one open() per entry, fanned out but not cached. Measured warm on
// 22 cores: 5000 entries take 29 ms serially and 6.9 ms across threads. The
// remaining ceiling is a cold cache or a network filesystem (NFS/SMB), where
// each entry costs a round trip. If that shows up: cache by (dev, ino, mtime),
// or skip the sniff for zero-byte files.
pub fn detect_kind(path: &Path, is_dir: bool) -> FileKind {
if is_dir {
return FileKind::Dir;
@@ -221,6 +310,9 @@ fn looks_like_text(head: &[u8]) -> bool {
}
}
/// Reported when a file's modification time is unavailable or unrepresentable.
const EPOCH_MTIME: &str = "1970-01-01T00:00:00Z";
fn mtime_str(m: &std::fs::Metadata) -> String {
let dt: Option<DateTime<chrono::Utc>> = m
.modified()
@@ -228,7 +320,7 @@ fn mtime_str(m: &std::fs::Metadata) -> String {
.and_then(|t| t.duration_since(UNIX_EPOCH).ok())
.and_then(|d| DateTime::from_timestamp(d.as_secs() as i64, 0));
dt.map(|d| d.to_rfc3339_opts(chrono::SecondsFormat::Secs, true))
.unwrap_or_else(|| "1970-01-01T00:00:00Z".to_string())
.unwrap_or_else(|| EPOCH_MTIME.to_string())
}
// ---------------------------------------------------------------------------
@@ -1190,6 +1282,84 @@ mod tests {
assert_eq!(detect_kind(&t.root.join("nope"), false), FileKind::Binary);
}
/// `SNIFF_BUDGET` is process-wide, so the two tests that assert on it must
/// not run at the same time as each other.
static BUDGET_TESTS: std::sync::Mutex<()> = std::sync::Mutex::new(());
/// A directory big enough to take the fan-out path must produce exactly
/// what the serial path would. Pass 2 fills `is_dir`, `size`, `mtime` and
/// `kind`, all on worker threads, so a chunk-boundary mistake would show up
/// as a row carrying another row's metadata or the untouched placeholders.
#[test]
fn parallel_rows_match_serial() {
use std::sync::atomic::Ordering;
let _guard = BUDGET_TESTS.lock().unwrap();
let t = T::new();
let big = t.root.join("big");
std::fs::create_dir_all(&big).unwrap();
// Well over SNIFF_CHUNK, so several chunks are handed out.
let n = SNIFF_CHUNK * 3 + 7;
for i in 0..n {
let p = big.join(format!("f{i:05}"));
// Every file gets a distinct length, so `size` pins the row identity.
let pad = vec![b'A'; i];
let mut body = match i % 3 {
0 => vec![0x89, b'P', b'N', b'G', 0x0D, 0x0A, 0x1A, 0x0A],
1 => b"plain text\n".to_vec(),
_ => vec![0u8, 1, 2, 3],
};
body.extend_from_slice(&pad);
std::fs::write(&p, &body).unwrap();
}
std::fs::create_dir(big.join("a-subdir")).unwrap();
let entries = list_dir(&big).unwrap();
assert_eq!(entries.len(), n + 1);
// Every row must agree with what a single-threaded read of that same
// file reports: kind, size and directory flag.
for e in &entries {
let p = big.join(&e.name);
let m = std::fs::metadata(&p).unwrap();
assert_eq!(e.is_dir, m.is_dir(), "{}", e.name);
assert_eq!(e.size, m.len(), "size mismatch on {}", e.name);
assert_eq!(e.mtime, mtime_str(&m), "mtime mismatch on {}", e.name);
assert_eq!(
e.kind,
detect_kind(&p, e.is_dir),
"kind mismatch on {}",
e.name
);
}
// Sanity: several kinds are actually present, so the loop above is not
// trivially true, and the directory sorts first.
let kinds: Vec<FileKind> = entries.iter().map(|e| e.kind).collect();
assert!(kinds.contains(&FileKind::Image));
assert!(kinds.contains(&FileKind::Text));
assert!(kinds.contains(&FileKind::Binary));
assert_eq!(entries[0].name, "a-subdir");
assert_eq!(entries[0].kind, FileKind::Dir);
// The thread budget is fully returned once the listing is done.
assert_eq!(SNIFF_BUDGET.load(Ordering::Acquire), 0);
}
#[test]
fn sniff_permit_never_exceeds_the_budget() {
use std::sync::atomic::Ordering;
let _guard = BUDGET_TESTS.lock().unwrap();
let total = sniff_budget_total();
let a = SniffPermit::claim(total * 2);
assert_eq!(a.0, total, "a single claim is capped at the total");
// Nothing left: the next listing runs serially rather than queueing.
let b = SniffPermit::claim(4);
assert_eq!(b.0, 0);
drop(a);
drop(b);
assert_eq!(SNIFF_BUDGET.load(Ordering::Acquire), 0);
}
#[test]
fn list_dir_reports_kinds() {
let t = T::new();