From 47c60788ff1596c3b2d2995a72c5698f11a8d6be Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fran=C3=A7ois-Xavier=20Talbot?= <108630700+fetchfern@users.noreply.github.com> Date: Wed, 24 Jun 2026 00:19:10 -0400 Subject: [PATCH] attributions fixes (#6487) * feat(labrinth): begin & commit one transaction per file scans * chore(labrinth): use ro pool for the query --- apps/labrinth/src/auth/checks.rs | 9 ++-- apps/labrinth/src/lib.rs | 3 +- apps/labrinth/src/queue/file_scan.rs | 52 +++++++++++---------- apps/labrinth/src/queue/moderation.rs | 10 +++- apps/labrinth/src/routes/updates.rs | 4 +- apps/labrinth/src/routes/v2/projects.rs | 4 +- apps/labrinth/src/routes/v2/versions.rs | 27 ++++++++--- apps/labrinth/src/routes/v3/projects.rs | 8 +++- apps/labrinth/src/routes/v3/version_file.rs | 1 + apps/labrinth/src/routes/v3/versions.rs | 46 +++++++++++++----- 10 files changed, 110 insertions(+), 54 deletions(-) diff --git a/apps/labrinth/src/auth/checks.rs b/apps/labrinth/src/auth/checks.rs index 970f9ab7b2..29509239f4 100644 --- a/apps/labrinth/src/auth/checks.rs +++ b/apps/labrinth/src/auth/checks.rs @@ -1,10 +1,10 @@ use crate::database; -use crate::database::PgPool; use crate::database::models::project_item::ProjectQueryResult; use crate::database::models::version_item::VersionQueryResult; use crate::database::models::{DBCollection, DBOrganization, DBTeamMember}; use crate::database::redis::RedisPool; use crate::database::{DBProject, DBVersion, models}; +use crate::database::{PgPool, ReadOnlyPgPool}; use crate::models::ids::FileId; use crate::models::projects::{ MissingAttributionFile, OverrideSource, Version, @@ -19,10 +19,10 @@ use itertools::Itertools; pub async fn enrich_dependency_attributions( versions: &mut [VersionQueryResult], - pool: &PgPool, + pool: &ReadOnlyPgPool, ) { let version_ids = versions.iter().map(|v| v.inner.id).collect::>(); - let dep_attr = get_dependency_attributions(pool, &version_ids) + let dep_attr = get_dependency_attributions(&**pool, &version_ids) .await .unwrap_or_default(); @@ -222,6 +222,7 @@ pub async fn filter_visible_versions( mut versions: Vec, user_option: &Option, pool: &PgPool, + ro_pool: &ReadOnlyPgPool, redis: &RedisPool, ) -> Result, ApiError> { let filtered_version_ids = filter_visible_version_ids( @@ -238,7 +239,7 @@ pub async fn filter_visible_versions( .await .unwrap_or_default(); - enrich_dependency_attributions(&mut versions, pool).await; + enrich_dependency_attributions(&mut versions, ro_pool).await; Ok(versions .into_iter() diff --git a/apps/labrinth/src/lib.rs b/apps/labrinth/src/lib.rs index 1d64722707..20e09e8c1d 100644 --- a/apps/labrinth/src/lib.rs +++ b/apps/labrinth/src/lib.rs @@ -101,10 +101,11 @@ pub fn app_setup( { let automated_moderation_queue_ref = automated_moderation_queue.clone(); let pool_ref = pool.clone(); + let ro_pool_ref = ro_pool.clone(); let redis_pool_ref = redis_pool.clone(); actix_rt::spawn(async move { automated_moderation_queue_ref - .task(pool_ref, redis_pool_ref) + .task(pool_ref, ro_pool_ref, redis_pool_ref) .await; }); } diff --git a/apps/labrinth/src/queue/file_scan.rs b/apps/labrinth/src/queue/file_scan.rs index d3978c4022..f9af623eff 100644 --- a/apps/labrinth/src/queue/file_scan.rs +++ b/apps/labrinth/src/queue/file_scan.rs @@ -44,8 +44,6 @@ pub async fn scan_all_files( redis: &RedisPool, file_host: &dyn FileHost, ) -> Result<()> { - let mut txn = db.begin().await.wrap_err("beginning transaction")?; - let files_to_scan = sqlx::query!( r#" select @@ -59,13 +57,13 @@ pub async fn scan_all_files( where fa.attributions_scanned_at is null "# ) - .fetch_all(&mut txn) + .fetch_all(db) .await .wrap_err("fetching files to scan")?; info!("Found {} files to scan", files_to_scan.len()); - let mut scanned_ids = Vec::new(); + let mut scanned_count = 0; for row in files_to_scan { let human_file_id = FileId::from(row.file_id); @@ -74,6 +72,10 @@ 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, @@ -108,33 +110,33 @@ pub async fn scan_all_files( })?; } - scanned_ids.push(file_id.0); + let now = Utc::now(); + sqlx::query!( + " + update file_scans + set attributions_scanned_at = $2 + where file_id = $1 + ", + file_id.0, + now, + ) + .execute(&mut txn) + .await + .wrap_err("marking file as scanned")?; + + txn.commit() + .await + .wrap_err("committing file scan transaction")?; + eyre::Ok(()) } .instrument(span) .await?; + + scanned_count += 1; } - if !scanned_ids.is_empty() { - let now = Utc::now(); - sqlx::query!( - " - update file_scans - set attributions_scanned_at = now - from unnest($1::bigint[], $2::timestamptz[]) as u(id, now) - where file_scans.file_id = u.id - ", - &scanned_ids, - &vec![now; scanned_ids.len()], - ) - .execute(&mut txn) - .await - .wrap_err("marking files as scanned")?; - } - - info!("Marked {} files as scanned", scanned_ids.len()); - - txn.commit().await.wrap_err("committing transaction")?; + info!("Marked {} files as scanned", scanned_count); Ok(()) } diff --git a/apps/labrinth/src/queue/moderation.rs b/apps/labrinth/src/queue/moderation.rs index 48f55a6c06..9efb0c40fb 100644 --- a/apps/labrinth/src/queue/moderation.rs +++ b/apps/labrinth/src/queue/moderation.rs @@ -1,10 +1,10 @@ use crate::auth::checks::filter_visible_versions; use crate::database; -use crate::database::PgPool; use crate::database::models::DBUserId; use crate::database::models::notification_item::NotificationBuilder; use crate::database::models::thread_item::ThreadMessageBuilder; use crate::database::redis::RedisPool; +use crate::database::{PgPool, ReadOnlyPgPool}; use crate::env::ENV; use crate::models::ids::ProjectId; use crate::models::notifications::NotificationBody; @@ -229,7 +229,12 @@ impl Default for AutomatedModerationQueue { } impl AutomatedModerationQueue { - pub async fn task(&self, pool: PgPool, redis: RedisPool) { + pub async fn task( + &self, + pool: PgPool, + ro_pool: ReadOnlyPgPool, + redis: RedisPool, + ) { loop { let projects = self.projects.clone(); self.projects.clear(); @@ -375,6 +380,7 @@ impl AutomatedModerationQueue { .await?, &None, &pool, + &ro_pool, &redis, ) .await?; diff --git a/apps/labrinth/src/routes/updates.rs b/apps/labrinth/src/routes/updates.rs index 987d695534..17a4b5dc22 100644 --- a/apps/labrinth/src/routes/updates.rs +++ b/apps/labrinth/src/routes/updates.rs @@ -1,6 +1,6 @@ use std::collections::HashMap; -use crate::database::PgPool; +use crate::database::{PgPool, ReadOnlyPgPool}; use crate::env::ENV; use actix_web::{HttpRequest, HttpResponse, get, web}; use serde::{Deserialize, Serialize}; @@ -36,6 +36,7 @@ pub async fn forge_updates( web::Query(neo): web::Query, info: web::Path<(String,)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -82,6 +83,7 @@ pub async fn forge_updates( .collect(), &user_option, &pool, + &ro_pool, &redis, ) .await?; diff --git a/apps/labrinth/src/routes/v2/projects.rs b/apps/labrinth/src/routes/v2/projects.rs index 700001b5f3..ded7348367 100644 --- a/apps/labrinth/src/routes/v2/projects.rs +++ b/apps/labrinth/src/routes/v2/projects.rs @@ -1,7 +1,7 @@ -use crate::database::PgPool; use crate::database::models::categories::LinkPlatform; use crate::database::models::{project_item, version_item}; use crate::database::redis::RedisPool; +use crate::database::{PgPool, ReadOnlyPgPool}; use crate::file_hosting::FileHost; use crate::models::projects::{ Link, MonetizationStatus, Project, ProjectStatus, Version, @@ -366,6 +366,7 @@ pub async fn dependency_list( req: HttpRequest, info: web::Path<(String,)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -374,6 +375,7 @@ pub async fn dependency_list( req, info, pool.clone(), + ro_pool, redis.clone(), session_queue, ) diff --git a/apps/labrinth/src/routes/v2/versions.rs b/apps/labrinth/src/routes/v2/versions.rs index 6184ea2136..d444bf8ed2 100644 --- a/apps/labrinth/src/routes/v2/versions.rs +++ b/apps/labrinth/src/routes/v2/versions.rs @@ -1,8 +1,8 @@ use std::collections::HashMap; use super::ApiError; -use crate::database::PgPool; use crate::database::redis::RedisPool; +use crate::database::{PgPool, ReadOnlyPgPool}; use crate::models; use crate::models::ids::VersionId; use crate::models::projects::{ @@ -89,6 +89,7 @@ pub async fn version_list( info: web::Path<(String,)>, web::Query(filters): web::Query, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -147,6 +148,7 @@ pub async fn version_list( info, web::Query(filters), pool, + ro_pool, redis, session_queue, ) @@ -196,6 +198,7 @@ pub async fn version_project_get( req: HttpRequest, info: web::Path<(String, String)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -204,6 +207,7 @@ pub async fn version_project_get( req, id, pool, + ro_pool, redis, session_queue, ) @@ -238,6 +242,7 @@ pub async fn versions_get( req: HttpRequest, web::Query(ids): web::Query, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -249,6 +254,7 @@ pub async fn versions_get( req, web::Query(ids), pool, + ro_pool, redis, session_queue, ) @@ -286,15 +292,22 @@ pub async fn version_get( req: HttpRequest, info: web::Path<(models::ids::VersionId,)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { let id = info.into_inner().0; - let response = - v3::versions::version_get_helper(req, id, pool, redis, session_queue) - .await - .map(|b| HttpResponse::Ok().json(b)) - .or_else(v2_reroute::flatten_404_error)?; + let response = v3::versions::version_get_helper( + req, + id, + pool, + ro_pool, + redis, + session_queue, + ) + .await + .map(|b| HttpResponse::Ok().json(b)) + .or_else(v2_reroute::flatten_404_error)?; // Convert response to V2 format match v2_reroute::extract_ok_json::(response).await { Ok(version) => { @@ -364,6 +377,7 @@ pub async fn version_edit( req: HttpRequest, info: web::Path<(VersionId,)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, new_version: web::Json, session_queue: web::Data, @@ -384,6 +398,7 @@ pub async fn version_edit( req.clone(), (*info).0, pool.clone(), + ro_pool.clone(), redis.clone(), session_queue.clone(), ) diff --git a/apps/labrinth/src/routes/v3/projects.rs b/apps/labrinth/src/routes/v3/projects.rs index cfad6a004b..4d8007a48b 100644 --- a/apps/labrinth/src/routes/v3/projects.rs +++ b/apps/labrinth/src/routes/v3/projects.rs @@ -11,7 +11,7 @@ use crate::database::models::{ }; use crate::database::redis::RedisPool; use crate::database::{self, models as db_models}; -use crate::database::{PgPool, PgTransaction}; +use crate::database::{PgPool, PgTransaction, ReadOnlyPgPool}; use crate::env::ENV; use crate::file_hosting::{FileHost, FileHostPublicity}; use crate::models::ids::{ProjectId, VersionId}; @@ -1289,16 +1289,19 @@ pub async fn dependency_list( req: HttpRequest, info: web::Path<(String,)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { - dependency_list_internal(req, info, pool, redis, session_queue).await + dependency_list_internal(req, info, pool, ro_pool, redis, session_queue) + .await } pub async fn dependency_list_internal( req: HttpRequest, info: web::Path<(String,)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -1372,6 +1375,7 @@ pub async fn dependency_list_internal( versions_result, &user_option, &pool, + &ro_pool, &redis, ) .await?; diff --git a/apps/labrinth/src/routes/v3/version_file.rs b/apps/labrinth/src/routes/v3/version_file.rs index f182edaed3..ad6912103b 100644 --- a/apps/labrinth/src/routes/v3/version_file.rs +++ b/apps/labrinth/src/routes/v3/version_file.rs @@ -251,6 +251,7 @@ pub async fn get_versions_from_hashes( .await?, &user_option, &pool, + &pool, &redis, ) .await?; diff --git a/apps/labrinth/src/routes/v3/versions.rs b/apps/labrinth/src/routes/v3/versions.rs index 1de3408440..e8690fa83f 100644 --- a/apps/labrinth/src/routes/v3/versions.rs +++ b/apps/labrinth/src/routes/v3/versions.rs @@ -7,7 +7,6 @@ use crate::auth::checks::{ }; use crate::auth::get_user_from_headers; use crate::database; -use crate::database::PgPool; use crate::database::models::loader_fields::{ self, LoaderField, LoaderFieldEnumValue, VersionField, }; @@ -16,6 +15,7 @@ use crate::database::models::version_item::{ }; use crate::database::models::{DBOrganization, image_item}; use crate::database::redis::RedisPool; +use crate::database::{PgPool, ReadOnlyPgPool}; use crate::models; use crate::models::ids::VersionId; use crate::models::images::ImageContext; @@ -64,16 +64,19 @@ pub async fn version_project_get( req: HttpRequest, info: web::Path<(String, String)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { let info = info.into_inner(); - version_project_get_helper(req, info, pool, redis, session_queue).await + version_project_get_helper(req, info, pool, ro_pool, redis, session_queue) + .await } pub async fn version_project_get_helper( req: HttpRequest, id: (String, String), pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -118,7 +121,7 @@ pub async fn version_project_get_helper( let version_id = version.inner.id; enrich_dependency_attributions( std::slice::from_mut(&mut version), - &pool, + &ro_pool, ) .await; let mut v = models::projects::Version::from(version); @@ -168,6 +171,7 @@ pub async fn versions_get( req: HttpRequest, web::Query(ids): web::Query, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -191,9 +195,14 @@ pub async fn versions_get( .map(|x| x.1) .ok(); - let mut versions = - filter_visible_versions(versions_data, &user_option, &pool, &redis) - .await?; + let mut versions = filter_visible_versions( + versions_data, + &user_option, + &pool, + &ro_pool, + &redis, + ) + .await?; if !ids.include_changelog { for version in &mut versions { @@ -208,17 +217,19 @@ pub async fn version_get( req: HttpRequest, info: web::Path<(models::ids::VersionId,)>, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result, ApiError> { let id = info.into_inner().0; - version_get_helper(req, id, pool, redis, session_queue).await + version_get_helper(req, id, pool, ro_pool, redis, session_queue).await } pub async fn version_get_helper( req: HttpRequest, id: models::ids::VersionId, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result, ApiError> { @@ -240,8 +251,11 @@ pub async fn version_get_helper( && is_visible_version(&data.inner, &user_option, &pool, &redis).await? { let version_id = data.inner.id; - enrich_dependency_attributions(std::slice::from_mut(&mut data), &pool) - .await; + enrich_dependency_attributions( + std::slice::from_mut(&mut data), + &ro_pool, + ) + .await; let mut version = models::projects::Version::from(data); let missing = get_files_missing_attribution(&**pool, &[version_id]) .await @@ -797,6 +811,7 @@ async fn version_list( info: web::Path<(String,)>, web::Query(filters): web::Query, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -805,6 +820,7 @@ async fn version_list( info, web::Query(filters), pool, + ro_pool, redis, session_queue, ) @@ -816,6 +832,7 @@ pub async fn version_list_internal( info: web::Path<(String,)>, web::Query(filters): web::Query, pool: web::Data, + ro_pool: web::Data, redis: web::Data, session_queue: web::Data, ) -> Result { @@ -955,9 +972,14 @@ pub async fn version_list_internal( }); response.dedup_by(|a, b| a.inner.id == b.inner.id); - let mut response = - filter_visible_versions(response, &user_option, &pool, &redis) - .await?; + let mut response = filter_visible_versions( + response, + &user_option, + &pool, + &ro_pool, + &redis, + ) + .await?; if !filters.include_changelog { for version in &mut response {