fix(python): re-verify wheel RECORD on local cache reuse (once per worker) (#9775)

A corrupt local pip cache entry (e.g. wmill==1.739.0 missing s3_reader.py
after out-of-band file loss on a persistent/shared cache volume) was trusted
indefinitely: handle_python_reqs only checked the .valid.windmill marker on
the reuse fast path. verify_wheel_record already guarded the install and
S3-pull paths, but never ran again once the marker existed.

Re-verify the wheel RECORD on the first reuse of each cache entry per worker
process and repair (wipe + reinstall) on failure. A VERIFIED_VENVS in-memory
set makes every subsequent reuse skip the scan, so the warm-cache hot path
keeps paying only its original single stat.

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-06-25 13:58:49 +00:00
committed by GitHub
co-authored by Claude Opus 4.8
parent aa098c70c0
commit 6c71c33470
+140 -2
View File
@@ -2050,6 +2050,15 @@ Returned from server: py_version - {:?}, py_version_v2 - {:?}
lazy_static::lazy_static! {
static ref PIP_SECRET_VARIABLE: Regex = Regex::new(r"\$\{PIP_SECRET:([^\s\}]+)\}").unwrap();
/// venv paths whose wheel RECORD this process has already verified against
/// disk. A cache entry is only damaged out-of-band (disk-pressure eviction,
/// an interrupted extraction on a shared cache volume, or a corrupt entry
/// that predates this worker), never spontaneously while we keep running, so
/// re-verifying it once per process is enough — every later reuse trusts the
/// in-memory marker and pays only the original single stat.
static ref VERIFIED_VENVS: tokio::sync::Mutex<HashSet<String>> =
tokio::sync::Mutex::new(HashSet::new());
}
/// Spawn process of uv install
@@ -2489,8 +2498,53 @@ pub async fn handle_python_reqs(
req.replace(' ', "").replace('/', "").replace(':', "")
);
if metadata(venv_p.clone() + "/.valid.windmill").await.is_ok() {
req_paths.push(venv_p);
in_cache.push(req.to_string());
// The .valid.windmill marker is written once at creation time, after
// verify_wheel_record passes on the install/pull paths. It is an empty
// file with no binding to the directory contents, so a file dropped
// out-of-band afterwards (disk-pressure eviction, interrupted tar
// extraction on a shared cache volume, or a corrupt entry that
// predates this worker) leaves the marker intact while the wheel is
// incomplete. Re-verify the RECORD once per process so such an entry
// is repaired rather than trusted; VERIFIED_VENVS makes every later
// reuse skip the scan and pay only the single stat above.
let already_verified = VERIFIED_VENVS.lock().await.contains(&venv_p);
let verify_res = if already_verified {
Ok(())
} else {
verify_wheel_record(&venv_p).await
};
match verify_res {
Ok(()) => {
if !already_verified {
VERIFIED_VENVS.lock().await.insert(venv_p.clone());
}
req_paths.push(venv_p);
in_cache.push(req.to_string());
}
Err(verify_err) => {
tracing::warn!(
workspace_id = %w_id,
job_id = %job_id,
"Local cache for {venv_p} failed wheel RECORD verification, will reinstall: {verify_err}"
);
append_logs(
&job_id,
w_id,
format!(
"\n[!] cached wheel for {req} failed integrity check, reinstalling: {verify_err}\n"
),
conn,
)
.await;
if let Err(rm_err) = tokio::fs::remove_dir_all(&venv_p).await {
tracing::warn!(
workspace_id = %w_id,
"could not remove broken cache dir {venv_p}: {rm_err}"
);
}
req_with_penv.push((req.to_string(), venv_p));
}
}
} else {
// There is no valid or no wheel at all. Regardless of if there is content or not, we will overwrite it with --reinstall flag
req_with_penv.push((req.to_string(), venv_p));
@@ -3467,4 +3521,88 @@ mod tests {
assert_eq!(kept, lines(&["# py: 3.11", "requests==2.0"]));
assert_eq!(ignored, lines(&["pyyaml==6.0"]));
}
/// Materialize a fake installed wheel: every file in `files` is created, and
/// `record_entries` is written verbatim as the RECORD (so a test can list a
/// path in RECORD without creating it, to simulate out-of-band loss).
fn write_fake_wheel(root: &std::path::Path, files: &[&str], record_entries: &[&str]) {
for f in files {
let full = root.join(f);
std::fs::create_dir_all(full.parent().unwrap()).unwrap();
std::fs::write(full, b"x").unwrap();
}
let dist_info = root.join("pkg-1.0.0.dist-info");
std::fs::create_dir_all(&dist_info).unwrap();
std::fs::write(
dist_info.join("RECORD"),
record_entries.join("\n") + "\n",
)
.unwrap();
}
#[tokio::test]
async fn test_verify_wheel_record_ok_when_all_present() {
let dir = tempfile::tempdir().unwrap();
write_fake_wheel(
dir.path(),
&["pkg/__init__.py", "pkg/mod.py"],
&[
"pkg/__init__.py,sha256=aaa,1",
"pkg/mod.py,sha256=bbb,1",
"pkg-1.0.0.dist-info/RECORD,,",
],
);
assert!(verify_wheel_record(dir.path().to_str().unwrap())
.await
.is_ok());
}
#[tokio::test]
async fn test_verify_wheel_record_err_when_file_missing() {
let dir = tempfile::tempdir().unwrap();
// RECORD lists pkg/mod.py but we never create it: the exact failure mode
// the customer hit (wmill/s3_reader.py present in RECORD, gone on disk).
write_fake_wheel(
dir.path(),
&["pkg/__init__.py"],
&[
"pkg/__init__.py,sha256=aaa,1",
"pkg/mod.py,sha256=bbb,1",
"pkg-1.0.0.dist-info/RECORD,,",
],
);
let err = verify_wheel_record(dir.path().to_str().unwrap())
.await
.unwrap_err();
assert!(err.contains("pkg/mod.py"), "unexpected error: {err}");
}
#[tokio::test]
async fn test_verify_wheel_record_err_when_no_dist_info() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("loose.py"), b"x").unwrap();
assert!(verify_wheel_record(dir.path().to_str().unwrap())
.await
.is_err());
}
#[tokio::test]
async fn test_verify_wheel_record_skips_absolute_and_escaping_entries() {
let dir = tempfile::tempdir().unwrap();
// Absolute and `..` RECORD entries are not package-relative and must be
// skipped rather than reported missing.
write_fake_wheel(
dir.path(),
&["pkg/__init__.py"],
&[
"pkg/__init__.py,sha256=aaa,1",
"/etc/passwd,sha256=ccc,1",
"../outside.py,sha256=ddd,1",
"pkg-1.0.0.dist-info/RECORD,,",
],
);
assert!(verify_wheel_record(dir.path().to_str().unwrap())
.await
.is_ok());
}
}