Batch pending file scanning ops (#6491)

* Batch pending file scanning ops

* prepare

* move some stuff out of the txn

* scan concurrency

* prepr

* 🦀 automod attribution is dead, long live automod

* fix ci

* update scan file filtering
This commit is contained in:
aecsocket
2026-06-24 21:46:07 +00:00
committed by GitHub
parent 586cdc90f9
commit f4a2999b20
6 changed files with 176 additions and 394 deletions
+139 -34
View File
@@ -1,11 +1,12 @@
use std::collections::HashMap;
use std::io::{Cursor, Read};
use std::sync::Arc;
use chrono::Utc;
use eyre::{Result, eyre};
use hex::ToHex;
use sha1::Digest;
use tokio::task::spawn_blocking;
use tokio::task::{spawn, spawn_blocking};
use tracing::{Instrument, info, info_span, warn};
use zip::ZipArchive;
@@ -29,6 +30,15 @@ use crate::queue::moderation::{
use crate::util::error::Context;
use crate::util::http::HTTP_CLIENT;
const PENDING_FILE_SCAN_BATCH_SIZE: i64 = 100;
#[derive(Clone)]
struct PendingFileScan {
file_id: DBFileId,
url: String,
project_id: DBProjectId,
}
/// Attribution enforcement is version/project-scoped, not file-hash-scoped.
///
/// Versions or projects listed in `attributions_exemptions` predate this
@@ -39,11 +49,50 @@ use crate::util::http::HTTP_CLIENT;
/// versions must go through the `attribution_enforced_versions` view so
/// grandfathered versions and projects are ignored without making the SHA1
/// itself exempt.
pub async fn scan_all_files(
pub async fn scan_all_pending_files(
db: &PgPool,
redis: &RedisPool,
file_host: &dyn FileHost,
file_host: Arc<dyn FileHost>,
) -> Result<()> {
let scan_concurrency = ENV.FILE_SCAN_CONCURRENCY.max(1);
let total_to_scan = sqlx::query_scalar!(
r#"
select count(*) as "count!" from file_scans
where attributions_scanned_at is null
"#,
)
.fetch_one(db)
.await
.wrap_err("fetching number of files to scan")?;
info!(
"Found {total_to_scan} total pending files to scan, running in batches of {PENDING_FILE_SCAN_BATCH_SIZE} with concurrency {scan_concurrency}"
);
loop {
let scanned_count = scan_pending_files_batch(
db,
redis,
file_host.clone(),
scan_concurrency * PENDING_FILE_SCAN_BATCH_SIZE,
)
.await?;
if scanned_count == 0 {
break;
}
}
Ok(())
}
async fn scan_pending_files_batch(
db: &PgPool,
redis: &RedisPool,
file_host: Arc<dyn FileHost>,
scan_limit: i64,
) -> Result<usize> {
let files_to_scan = sqlx::query!(
r#"
select
@@ -55,14 +104,61 @@ pub async fn scan_all_files(
inner join attribution_enforced_versions aev on aev.id = f.version_id
inner join versions v on v.id = f.version_id
where fa.attributions_scanned_at is null
"#
order by fa.file_id
limit $1
"#,
scan_limit,
)
.fetch_all(db)
.await
.wrap_err("fetching files to scan")?;
info!("Found {} files to scan", files_to_scan.len());
info!(
"Found {} pending files to scan, splitting into jobs of {PENDING_FILE_SCAN_BATCH_SIZE}",
files_to_scan.len(),
);
let files_to_scan: Vec<_> = files_to_scan
.into_iter()
.map(|row| PendingFileScan {
file_id: row.file_id,
url: row.url,
project_id: row.project_id,
})
.collect();
let mut tasks = Vec::new();
for chunk in files_to_scan.chunks(PENDING_FILE_SCAN_BATCH_SIZE as usize) {
let db = db.clone();
let redis = redis.clone();
let file_host = file_host.clone();
let chunk = chunk.to_vec();
tasks.push(spawn(async move {
scan_pending_files_chunk(&db, &redis, &*file_host, chunk).await
}));
}
let mut scanned_count = 0;
for task in tasks {
scanned_count += task
.await
.wrap_err("joining file scan task")?
.wrap_err("scanning pending file chunk")?;
}
info!("Marked {} files as scanned", scanned_count);
Ok(scanned_count)
}
async fn scan_pending_files_chunk(
db: &PgPool,
redis: &RedisPool,
file_host: &dyn FileHost,
files_to_scan: Vec<PendingFileScan>,
) -> Result<usize> {
info!("Scanning {} files", files_to_scan.len());
let mut scanned_count = 0;
for row in files_to_scan {
@@ -72,10 +168,6 @@ pub async fn scan_all_files(
info!("Scanning file");
let file_id = row.file_id;
let mut txn = db
.begin()
.await
.wrap_err("beginning file scan transaction")?;
let overrides = extract_override_files_from_storage(
file_host, file_id, &row.url,
@@ -87,29 +179,35 @@ pub async fn scan_all_files(
if overrides.is_empty() {
info!("Found no overrides");
} else {
info!("Found {} overrides", overrides.len());
return Ok(());
}
let resolved = resolve_overrides(&overrides, redis, &mut txn)
.await
.wrap_err_with(|| {
eyre!("resolving overrides for file {file_id:?}")
})?;
info!("Resolved: {resolved:#?}");
info!("Found {} overrides", overrides.len());
persist_attribution_results(
row.project_id,
file_id,
&overrides,
&resolved,
redis,
&mut txn,
)
let mut txn = db
.begin()
.await
.wrap_err("beginning file scan transaction")?;
let resolved = resolve_overrides(&overrides, redis, &mut txn)
.await
.wrap_err_with(|| {
eyre!("persisting attribution results for file {file_id:?}")
eyre!("resolving overrides for file {file_id:?}")
})?;
}
info!("Resolved: {resolved:#?}");
persist_attribution_results(
row.project_id,
file_id,
&overrides,
&resolved,
redis,
&mut txn,
)
.await
.wrap_err_with(|| {
eyre!("persisting attribution results for file {file_id:?}")
})?;
let now = Utc::now();
sqlx::query!(
@@ -137,9 +235,7 @@ pub async fn scan_all_files(
scanned_count += 1;
}
info!("Marked {} files as scanned", scanned_count);
Ok(())
Ok(scanned_count)
}
pub async fn scan_file(
@@ -276,10 +372,19 @@ fn extract_override_files(data: &[u8]) -> Result<Vec<OverrideFile>> {
continue;
}
let should_scan_file = name.contains(".jar")
|| (name.contains(".zip") && !name.ends_with(".zip.txt"));
if name.matches('/').count() > 2 || !should_scan_file {
let should_skip = name.starts_with("mods/.connector/")
|| name.starts_with(".sable/natives/")
|| name.starts_with("local/crash_assistant/")
|| name.starts_with("mods/mcef-libraries/")
|| name.starts_with("mods/mcef-cache/")
|| name.starts_with("config/super_resolution/libraries/")
|| name.starts_with("config/Veinminer/update/")
|| name.starts_with("config/epicfight/native/")
|| name.starts_with("essential/")
|| name.ends_with(".rpo")
|| name.ends_with(".txt");
let should_scan = name.contains(".jar") || name.contains(".zip");
if should_scan && !should_skip {
continue;
}