From 28856cb0cbf0d8f1bf2af8c767135f9c72e1de33 Mon Sep 17 00:00:00 2001 From: Denis Cornehl Date: Tue, 11 Aug 2026 22:54:31 +0200 Subject: [PATCH] build queue: refactor & cache priority-finding for new releases --- ...2e90e944d9009d75ad527e7442e312be98ea9.json | 32 -- ...91b40f873618036c55c679e5854dbd4cbf91e.json | 15 - ...da4e0fcf9c7edd8485731e4e6db04d752774.json} | 8 +- ...7ae600eef5f90aca5655143e716490f7d68b2.json | 33 ++ ...fdf9c313f54cd346e509c9d131c705c216b0e.json | 15 - ...628ebb710041c7397b4727a670e76beb1e2a7.json | 40 -- ...9b39ece40715ce25c6966760978ad8e6c6e15.json | 16 - ...bb821decc862cf905aa21feff167d401befbf.json | 28 ++ ...2fac00ff241eabb2760d66a909ca6fd76b668.json | 16 + ...3c241ba2baaeaa1d0534170c525627b2f05d.json} | 4 +- ...a0dd32dd4d0cf45060c16762b322a4e3c7ab5.json | 14 - Cargo.lock | 2 + ...621d5a82ef6a5333a617b1d2fb8631ebe9b42.json | 26 -- .../bin/docs_rs_watcher/src/index_watcher.rs | 15 +- crates/lib/docs_rs_build_queue/Cargo.toml | 1 + crates/lib/docs_rs_build_queue/src/lib.rs | 2 + .../lib/docs_rs_build_queue/src/priority.rs | 396 ++++++++++++++-- .../docs_rs_build_queue/src/queue/blocking.rs | 141 ++---- .../src/queue/non_blocking.rs | 423 +++++++----------- .../docs_rs_build_queue/src/testing/mod.rs | 1 + .../src/testing/test_env.rs | 117 +++++ .../src/workspaces.rs | 166 +++---- crates/lib/docs_rs_test_fakes/Cargo.toml | 1 + .../docs_rs_test_fakes/src/github_stats.rs | 78 ++++ crates/lib/docs_rs_test_fakes/src/legacy.rs | 59 +-- crates/lib/docs_rs_test_fakes/src/lib.rs | 4 +- 26 files changed, 932 insertions(+), 721 deletions(-) delete mode 100644 .sqlx/query-1002ada46a8b06269d7aa42acc52e90e944d9009d75ad527e7442e312be98ea9.json delete mode 100644 .sqlx/query-1fcb4371909a27b059e3dba7c0791b40f873618036c55c679e5854dbd4cbf91e.json rename .sqlx/{query-53eeab98d127a1ab001d09cc2288b8c185248d428b73c6519f7ea17557c8f02a.json => query-78749956731e17cee4cf54fa5284da4e0fcf9c7edd8485731e4e6db04d752774.json} (80%) create mode 100644 .sqlx/query-7ae6216099898467ff3e831b88d7ae600eef5f90aca5655143e716490f7d68b2.json delete mode 100644 .sqlx/query-7bea41567a747d6043891a7934cfdf9c313f54cd346e509c9d131c705c216b0e.json delete mode 100644 .sqlx/query-8a31ef6960f9769024fc7818761628ebb710041c7397b4727a670e76beb1e2a7.json delete mode 100644 .sqlx/query-a56472d4dbe0c773365671822ed9b39ece40715ce25c6966760978ad8e6c6e15.json create mode 100644 .sqlx/query-af1b135f6b08d1b7a8a4c79e1b3bb821decc862cf905aa21feff167d401befbf.json create mode 100644 .sqlx/query-b794f125a36f4ae162abe248f152fac00ff241eabb2760d66a909ca6fd76b668.json rename .sqlx/{query-926a6602b17498d10b30787583ad7f4082ee4fc478a6008dfc8dcae110d651fa.json => query-cdf712b5e258465fcdb6a9b132243c241ba2baaeaa1d0534170c525627b2f05d.json} (64%) delete mode 100644 .sqlx/query-e670555bc281dd478519f46656ca0dd32dd4d0cf45060c16762b322a4e3c7ab5.json delete mode 100644 crates/bin/docs_rs_builder/.sqlx/query-aad790aa7ef85357e7f57c83b21621d5a82ef6a5333a617b1d2fb8631ebe9b42.json create mode 100644 crates/lib/docs_rs_build_queue/src/testing/mod.rs create mode 100644 crates/lib/docs_rs_build_queue/src/testing/test_env.rs create mode 100644 crates/lib/docs_rs_test_fakes/src/github_stats.rs diff --git a/.sqlx/query-1002ada46a8b06269d7aa42acc52e90e944d9009d75ad527e7442e312be98ea9.json b/.sqlx/query-1002ada46a8b06269d7aa42acc52e90e944d9009d75ad527e7442e312be98ea9.json deleted file mode 100644 index 90c4d9e6aa..0000000000 --- a/.sqlx/query-1002ada46a8b06269d7aa42acc52e90e944d9009d75ad527e7442e312be98ea9.json +++ /dev/null @@ -1,32 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "INSERT INTO repositories (host, host_id, name, description, last_commit, stars, forks, issues, updated_at)\n VALUES ('github.com', $1, $2, 'Fake description!', NOW(), $3, $4, $5, NOW())\n RETURNING id", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id", - "type_info": "Int4", - "origin": { - "Table": { - "table": "repositories", - "name": "id" - } - } - } - ], - "parameters": { - "Left": [ - "Varchar", - "Varchar", - "Int4", - "Int4", - "Int4" - ] - }, - "nullable": [ - false - ] - }, - "hash": "1002ada46a8b06269d7aa42acc52e90e944d9009d75ad527e7442e312be98ea9" -} diff --git a/.sqlx/query-1fcb4371909a27b059e3dba7c0791b40f873618036c55c679e5854dbd4cbf91e.json b/.sqlx/query-1fcb4371909a27b059e3dba7c0791b40f873618036c55c679e5854dbd4cbf91e.json deleted file mode 100644 index 965eaddb33..0000000000 --- a/.sqlx/query-1fcb4371909a27b059e3dba7c0791b40f873618036c55c679e5854dbd4cbf91e.json +++ /dev/null @@ -1,15 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE queue\n SET priority = $2\n WHERE name = $1\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Text", - "Int4" - ] - }, - "nullable": [] - }, - "hash": "1fcb4371909a27b059e3dba7c0791b40f873618036c55c679e5854dbd4cbf91e" -} diff --git a/.sqlx/query-53eeab98d127a1ab001d09cc2288b8c185248d428b73c6519f7ea17557c8f02a.json b/.sqlx/query-78749956731e17cee4cf54fa5284da4e0fcf9c7edd8485731e4e6db04d752774.json similarity index 80% rename from .sqlx/query-53eeab98d127a1ab001d09cc2288b8c185248d428b73c6519f7ea17557c8f02a.json rename to .sqlx/query-78749956731e17cee4cf54fa5284da4e0fcf9c7edd8485731e4e6db04d752774.json index 07ac8b0e41..5a18df58af 100644 --- a/.sqlx/query-53eeab98d127a1ab001d09cc2288b8c185248d428b73c6519f7ea17557c8f02a.json +++ b/.sqlx/query-78749956731e17cee4cf54fa5284da4e0fcf9c7edd8485731e4e6db04d752774.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT\n c.name as \"name: KrateName\",\n repo.override_build_priority as \"priority!\"\n\n FROM\n crates AS c\n INNER JOIN releases AS r ON c.latest_version_id = r.id\n INNER JOIN repositories AS repo ON r.repository_id = repo.id\n\n WHERE\n c.name = ANY($1) AND\n repo.override_build_priority IS NOT NULL\n ", + "query": "\n SELECT\n c.name as \"name: KrateName\",\n repo.override_build_priority as \"priority!\"\n\n FROM\n crates AS c\n INNER JOIN releases AS r ON c.latest_version_id = r.id\n INNER JOIN repositories AS repo ON r.repository_id = repo.id\n\n WHERE\n repo.override_build_priority IS NOT NULL\n ", "describe": { "columns": [ { @@ -27,14 +27,12 @@ } ], "parameters": { - "Left": [ - "TextArray" - ] + "Left": [] }, "nullable": [ false, true ] }, - "hash": "53eeab98d127a1ab001d09cc2288b8c185248d428b73c6519f7ea17557c8f02a" + "hash": "78749956731e17cee4cf54fa5284da4e0fcf9c7edd8485731e4e6db04d752774" } diff --git a/.sqlx/query-7ae6216099898467ff3e831b88d7ae600eef5f90aca5655143e716490f7d68b2.json b/.sqlx/query-7ae6216099898467ff3e831b88d7ae600eef5f90aca5655143e716490f7d68b2.json new file mode 100644 index 0000000000..e42ee67546 --- /dev/null +++ b/.sqlx/query-7ae6216099898467ff3e831b88d7ae600eef5f90aca5655143e716490f7d68b2.json @@ -0,0 +1,33 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO repositories (\n host,\n host_id,\n name,\n description,\n last_commit,\n stars,\n forks,\n issues,\n updated_at,\n override_build_priority\n )\n VALUES (\n 'github.com',\n $1,\n $2,\n 'Fake description!',\n NOW(),\n $3,\n $4,\n $5,\n NOW(),\n $6\n )\n RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Int4", + "origin": { + "Table": { + "table": "repositories", + "name": "id" + } + } + } + ], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Int4", + "Int4", + "Int4", + "Int4" + ] + }, + "nullable": [ + false + ] + }, + "hash": "7ae6216099898467ff3e831b88d7ae600eef5f90aca5655143e716490f7d68b2" +} diff --git a/.sqlx/query-7bea41567a747d6043891a7934cfdf9c313f54cd346e509c9d131c705c216b0e.json b/.sqlx/query-7bea41567a747d6043891a7934cfdf9c313f54cd346e509c9d131c705c216b0e.json deleted file mode 100644 index 51d73e6fed..0000000000 --- a/.sqlx/query-7bea41567a747d6043891a7934cfdf9c313f54cd346e509c9d131c705c216b0e.json +++ /dev/null @@ -1,15 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE releases SET repository_id = $1 WHERE crate_id = (SELECT id FROM crates WHERE name = $2)", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Int4", - "Text" - ] - }, - "nullable": [] - }, - "hash": "7bea41567a747d6043891a7934cfdf9c313f54cd346e509c9d131c705c216b0e" -} diff --git a/.sqlx/query-8a31ef6960f9769024fc7818761628ebb710041c7397b4727a670e76beb1e2a7.json b/.sqlx/query-8a31ef6960f9769024fc7818761628ebb710041c7397b4727a670e76beb1e2a7.json deleted file mode 100644 index 762ecaa59c..0000000000 --- a/.sqlx/query-8a31ef6960f9769024fc7818761628ebb710041c7397b4727a670e76beb1e2a7.json +++ /dev/null @@ -1,40 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n SELECT\n c.name as \"name: KrateName\",\n repo.crate_count\n\n FROM\n crates AS c\n INNER JOIN releases AS r ON c.latest_version_id = r.id\n INNER JOIN repositories AS repo ON r.repository_id = repo.id\n\n WHERE c.name = ANY($1)\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "name: KrateName", - "type_info": "Text", - "origin": { - "Table": { - "table": "crates", - "name": "name" - } - } - }, - { - "ordinal": 1, - "name": "crate_count", - "type_info": "Int4", - "origin": { - "Table": { - "table": "repositories", - "name": "crate_count" - } - } - } - ], - "parameters": { - "Left": [ - "TextArray" - ] - }, - "nullable": [ - false, - false - ] - }, - "hash": "8a31ef6960f9769024fc7818761628ebb710041c7397b4727a670e76beb1e2a7" -} diff --git a/.sqlx/query-a56472d4dbe0c773365671822ed9b39ece40715ce25c6966760978ad8e6c6e15.json b/.sqlx/query-a56472d4dbe0c773365671822ed9b39ece40715ce25c6966760978ad8e6c6e15.json deleted file mode 100644 index 1a3d824d5e..0000000000 --- a/.sqlx/query-a56472d4dbe0c773365671822ed9b39ece40715ce25c6966760978ad8e6c6e15.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE queue\n SET priority = $3\n WHERE\n name = ANY($1) AND\n priority = $2\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "TextArray", - "Int4", - "Int4" - ] - }, - "nullable": [] - }, - "hash": "a56472d4dbe0c773365671822ed9b39ece40715ce25c6966760978ad8e6c6e15" -} diff --git a/.sqlx/query-af1b135f6b08d1b7a8a4c79e1b3bb821decc862cf905aa21feff167d401befbf.json b/.sqlx/query-af1b135f6b08d1b7a8a4c79e1b3bb821decc862cf905aa21feff167d401befbf.json new file mode 100644 index 0000000000..594d945b15 --- /dev/null +++ b/.sqlx/query-af1b135f6b08d1b7a8a4c79e1b3bb821decc862cf905aa21feff167d401befbf.json @@ -0,0 +1,28 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT\n c.name as \"name: KrateName\"\n\n FROM\n crates AS c\n INNER JOIN releases AS r ON c.latest_version_id = r.id\n INNER JOIN repositories AS repo ON r.repository_id = repo.id\n\n WHERE\n repo.crate_count >= $1\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "name: KrateName", + "type_info": "Text", + "origin": { + "Table": { + "table": "crates", + "name": "name" + } + } + } + ], + "parameters": { + "Left": [ + "Int4" + ] + }, + "nullable": [ + false + ] + }, + "hash": "af1b135f6b08d1b7a8a4c79e1b3bb821decc862cf905aa21feff167d401befbf" +} diff --git a/.sqlx/query-b794f125a36f4ae162abe248f152fac00ff241eabb2760d66a909ca6fd76b668.json b/.sqlx/query-b794f125a36f4ae162abe248f152fac00ff241eabb2760d66a909ca6fd76b668.json new file mode 100644 index 0000000000..1668b5528f --- /dev/null +++ b/.sqlx/query-b794f125a36f4ae162abe248f152fac00ff241eabb2760d66a909ca6fd76b668.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "\n UPDATE queue\n SET priority = updates.priority\n FROM UNNEST($1::text[], $2::int[]) AS updates(name, priority)\n WHERE\n queue.name = updates.name\n AND queue.priority = $3\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "TextArray", + "Int4Array", + "Int4" + ] + }, + "nullable": [] + }, + "hash": "b794f125a36f4ae162abe248f152fac00ff241eabb2760d66a909ca6fd76b668" +} diff --git a/.sqlx/query-926a6602b17498d10b30787583ad7f4082ee4fc478a6008dfc8dcae110d651fa.json b/.sqlx/query-cdf712b5e258465fcdb6a9b132243c241ba2baaeaa1d0534170c525627b2f05d.json similarity index 64% rename from .sqlx/query-926a6602b17498d10b30787583ad7f4082ee4fc478a6008dfc8dcae110d651fa.json rename to .sqlx/query-cdf712b5e258465fcdb6a9b132243c241ba2baaeaa1d0534170c525627b2f05d.json index 77a9952a13..53e07d6439 100644 --- a/.sqlx/query-926a6602b17498d10b30787583ad7f4082ee4fc478a6008dfc8dcae110d651fa.json +++ b/.sqlx/query-cdf712b5e258465fcdb6a9b132243c241ba2baaeaa1d0534170c525627b2f05d.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT DISTINCT name AS \"name: KrateName\"\n FROM queue\n WHERE priority = $1", + "query": "\n SELECT DISTINCT name AS \"name: KrateName\"\n FROM queue\n WHERE priority = $1\n ", "describe": { "columns": [ { @@ -24,5 +24,5 @@ false ] }, - "hash": "926a6602b17498d10b30787583ad7f4082ee4fc478a6008dfc8dcae110d651fa" + "hash": "cdf712b5e258465fcdb6a9b132243c241ba2baaeaa1d0534170c525627b2f05d" } diff --git a/.sqlx/query-e670555bc281dd478519f46656ca0dd32dd4d0cf45060c16762b322a4e3c7ab5.json b/.sqlx/query-e670555bc281dd478519f46656ca0dd32dd4d0cf45060c16762b322a4e3c7ab5.json deleted file mode 100644 index eb92e2efce..0000000000 --- a/.sqlx/query-e670555bc281dd478519f46656ca0dd32dd4d0cf45060c16762b322a4e3c7ab5.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE releases SET repository_id = $1", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Int4" - ] - }, - "nullable": [] - }, - "hash": "e670555bc281dd478519f46656ca0dd32dd4d0cf45060c16762b322a4e3c7ab5" -} diff --git a/Cargo.lock b/Cargo.lock index 5fad3b8d42..67e751bee2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2048,6 +2048,7 @@ dependencies = [ "futures-util", "opentelemetry", "pretty_assertions", + "regex", "sqlx", "test-case", "tokio", @@ -2418,6 +2419,7 @@ version = "0.1.0" dependencies = [ "anyhow", "base64 0.23.1", + "bon", "chrono", "docs_rs_cargo_metadata", "docs_rs_database", diff --git a/crates/bin/docs_rs_builder/.sqlx/query-aad790aa7ef85357e7f57c83b21621d5a82ef6a5333a617b1d2fb8631ebe9b42.json b/crates/bin/docs_rs_builder/.sqlx/query-aad790aa7ef85357e7f57c83b21621d5a82ef6a5333a617b1d2fb8631ebe9b42.json deleted file mode 100644 index d6ea915434..0000000000 --- a/crates/bin/docs_rs_builder/.sqlx/query-aad790aa7ef85357e7f57c83b21621d5a82ef6a5333a617b1d2fb8631ebe9b42.json +++ /dev/null @@ -1,26 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT cov.total_items\n FROM\n crates as c\n INNER JOIN releases AS r ON c.id = r.crate_id\n LEFT OUTER JOIN doc_coverage AS cov ON r.id = cov.release_id\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "total_items", - "type_info": "Int4", - "origin": { - "Table": { - "table": "doc_coverage", - "name": "total_items" - } - } - } - ], - "parameters": { - "Left": [] - }, - "nullable": [ - true - ] - }, - "hash": "aad790aa7ef85357e7f57c83b21621d5a82ef6a5333a617b1d2fb8631ebe9b42" -} diff --git a/crates/bin/docs_rs_watcher/src/index_watcher.rs b/crates/bin/docs_rs_watcher/src/index_watcher.rs index 68b964d79e..3d03d3ce94 100644 --- a/crates/bin/docs_rs_watcher/src/index_watcher.rs +++ b/crates/bin/docs_rs_watcher/src/index_watcher.rs @@ -5,7 +5,7 @@ use crate::{ }; use anyhow::{Context as _, Result}; use crates_index_diff::Change; -use docs_rs_build_queue::{PRIORITY_MANUAL_FROM_CRATES_IO, priority::get_crate_priority}; +use docs_rs_build_queue::PRIORITY_MANUAL_FROM_CRATES_IO; use docs_rs_context::Context; use docs_rs_database::{ crate_details::update_latest_version_id, @@ -117,8 +117,8 @@ pub(crate) async fn get_new_crates( let crates_added = process_changes(context, &changes, config).await; - if let Err(err) = context.build_queue()?.deprioritize_workspaces().await { - error!(?err, "error deprioritizing workspaces"); + if let Err(err) = context.build_queue()?.reevaluate_priorities().await { + error!(?err, "error reevaluating queued release priorities"); } // set the reference in the database @@ -186,10 +186,11 @@ async fn process_version_yank_status(context: &Context, release: &CrateVersion) } async fn process_version_added(context: &Context, release: &CrateVersion) -> Result<()> { - let mut conn = context.pool()?.get_async().await?; - let priority = get_crate_priority(&mut conn, &release.name).await?; - context - .build_queue()? + let build_queue = context.build_queue()?; + + let priority = build_queue.find_priority(&release.name).await?; + + build_queue .add_crate(&release.name, &release.version, priority) .await .with_context(|| { diff --git a/crates/lib/docs_rs_build_queue/Cargo.toml b/crates/lib/docs_rs_build_queue/Cargo.toml index f1ec57d57b..2910040178 100644 --- a/crates/lib/docs_rs_build_queue/Cargo.toml +++ b/crates/lib/docs_rs_build_queue/Cargo.toml @@ -25,6 +25,7 @@ docs_rs_uri = { path = "../docs_rs_uri" } docs_rs_utils = { path = "../docs_rs_utils" } futures-util = { workspace = true } opentelemetry = { workspace = true } +regex = { workspace = true } sqlx = { workspace = true } tokio = { workspace = true } tracing = { workspace = true } diff --git a/crates/lib/docs_rs_build_queue/src/lib.rs b/crates/lib/docs_rs_build_queue/src/lib.rs index 471cce3960..0be40b89f6 100644 --- a/crates/lib/docs_rs_build_queue/src/lib.rs +++ b/crates/lib/docs_rs_build_queue/src/lib.rs @@ -2,6 +2,8 @@ mod config; mod metrics; pub mod priority; mod queue; +#[cfg(test)] +pub(crate) mod testing; mod types; pub use config::Config; diff --git a/crates/lib/docs_rs_build_queue/src/priority.rs b/crates/lib/docs_rs_build_queue/src/priority.rs index 413d9b2574..44707bcb73 100644 --- a/crates/lib/docs_rs_build_queue/src/priority.rs +++ b/crates/lib/docs_rs_build_queue/src/priority.rs @@ -1,7 +1,172 @@ -use crate::PRIORITY_DEFAULT; -use anyhow::Result; -use docs_rs_types::KrateName; +use crate::{PRIORITY_DEFAULT, PRIORITY_DEPRIORITIZED}; +use anyhow::{Context as _, Result}; +use docs_rs_database::Pool; +use docs_rs_repository_stats::workspaces::{ + get_crate_names_for_repository_build_priorities, get_crate_names_from_big_workspaces, +}; +use docs_rs_types::{Duration, KrateName}; use futures_util::stream::TryStreamExt; +use regex::Regex; +use std::{ + collections::{HashMap, HashSet}, + time::Instant, +}; +use tokio::sync::Mutex; +use tracing::info; + +const PRIORITY_RELOAD_FREQUENCY: Duration = Duration::from_secs(300); // 5 minutes + +/// cached crate priorities. +/// +/// Load & caches necessary data to figure out the wanted priority for a crate build. +#[derive(Debug)] +pub(crate) struct PrioritiesCache { + inner: Mutex, +} + +#[derive(Debug)] +struct PrioritiesCacheInner { + last_reload: Option, + pool: Pool, + deprioritize_workspace_size: u16, + patterns: Vec<(Regex, i32)>, + workspace_overrides: HashMap, + big_workspace_crate_names: HashSet, +} + +impl PrioritiesCacheInner { + async fn reload(&mut self) -> Result<()> { + let mut conn = self.pool.get_async().await?; + + let patterns = list_crate_priorities(&mut conn) + .await? + .into_iter() + .map(|(pattern, priority)| -> Result<_> { + let re = compile_like_pattern(&pattern) + .with_context(|| format!("can't compile pattern {} into regex", pattern))?; + Ok((re, priority)) + }) + .collect::>>()?; + + let workspace_overrides = + get_crate_names_for_repository_build_priorities(&mut conn).await?; + let workspace_crate_names = + get_crate_names_from_big_workspaces(&mut conn, self.deprioritize_workspace_size) + .await?; + + self.patterns = patterns; + self.workspace_overrides = workspace_overrides; + self.big_workspace_crate_names = workspace_crate_names; + self.last_reload = Some(Instant::now()); + + info!( + patterns_len = self.patterns.len(), + workspace_overrides_len = self.workspace_overrides.len(), + workspace_crate_names_len = self.big_workspace_crate_names.len(), + deprioritize_workspace_size = self.deprioritize_workspace_size, + "loaded crate priorities" + ); + + Ok(()) + } + + fn priority_from_pattern(&self, krate: &KrateName) -> Option { + self.patterns + .iter() + .find_map(|(regex, prio)| regex.is_match(krate.as_str()).then_some(*prio)) + } + + fn priority_from_workspace_override(&self, krate: &KrateName) -> Option { + self.workspace_overrides.get(krate).copied() + } + + fn priority_from_big_workspaces(&self, krate: &KrateName) -> Option { + self.big_workspace_crate_names + .contains(krate) + .then_some(PRIORITY_DEPRIORITIZED) + } +} + +impl PrioritiesCache { + pub(crate) fn new(pool: Pool, deprioritize_workspace_size: u16) -> Self { + let inner = PrioritiesCacheInner { + pool, + last_reload: None, + deprioritize_workspace_size, + patterns: Vec::new(), + workspace_overrides: HashMap::new(), + big_workspace_crate_names: HashSet::new(), + }; + + PrioritiesCache { + inner: Mutex::new(inner), + } + } + + /// get the priority for a crate. + /// + /// Checks in order: + /// 1. priority overrides via name pattern + /// 2. workspace / repo priority overrides + /// 3. big workspace deprio + /// 4. default prio + pub(crate) async fn get(&self, krate: &KrateName) -> Result { + let mut inner = self.inner.lock().await; + + if inner + .last_reload + .is_none_or(|last_reload| last_reload.elapsed() > (*PRIORITY_RELOAD_FREQUENCY)) + { + inner.reload().await?; + }; + + Ok(inner + .priority_from_pattern(krate) + .or_else(|| inner.priority_from_workspace_override(krate)) + .or_else(|| inner.priority_from_big_workspaces(krate)) + .unwrap_or(PRIORITY_DEFAULT)) + } + + /// force a reload of the cached priority data, ignoring the reload frequency. + #[cfg(test)] + pub(crate) async fn reload(&self) -> Result<()> { + let mut inner = self.inner.lock().await; + inner.reload().await + } +} + +/// compile a postgres LIKE pattern to a regex so we can match it in rust. +/// +/// For now we compile the postgres pattern to regex, so we can easily revert this PR. +/// Later we can just write regexes into the table and directly use them. +fn compile_like_pattern(pattern: &str) -> Result { + let mut regex = String::from("^"); + let mut chars = pattern.chars(); + + while let Some(ch) = chars.next() { + match ch { + '%' => regex.push_str(".*"), + '_' => regex.push('.'), + + // Postgres LIKE uses backslash as the default escape char. + // So `\%` means literal percent, `\_` means literal underscore. + '\\' => { + if let Some(escaped) = chars.next() { + regex.push_str(®ex::escape(&escaped.to_string())); + } else { + regex.push_str(®ex::escape("\\")); + } + } + + literal => { + regex.push_str(®ex::escape(&literal.to_string())); + } + } + } + + regex.push('$'); + Ok(Regex::new(®ex)?) +} /// Get the build queue priority for a crate, returns the matching pattern too pub async fn list_crate_priorities(conn: &mut sqlx::PgConnection) -> Result> { @@ -29,13 +194,6 @@ pub async fn get_crate_pattern_and_priority( .map(|row| (row.pattern, row.priority))) } -/// Get the build queue priority for a crate -pub async fn get_crate_priority(conn: &mut sqlx::PgConnection, name: &KrateName) -> Result { - Ok(get_crate_pattern_and_priority(conn, name) - .await? - .map_or(PRIORITY_DEFAULT, |(_, priority)| priority)) -} - /// Set all crates that match [`pattern`] to have a certain priority /// /// Note: `pattern` is used in a `LIKE` statement, so it must follow the postgres like syntax @@ -74,11 +232,86 @@ pub async fn remove_crate_priority( #[cfg(test)] mod tests { use super::*; - use docs_rs_config::AppConfig as _; - use docs_rs_database::{Config, testing::TestDatabase}; - use docs_rs_opentelemetry::testing::TestMetrics; + use crate::testing::test_env::TestEnv; + use docs_rs_repository_stats::workspaces::rewrite_repository_stats; + use docs_rs_test_fakes::FakeGithubStats; + use docs_rs_types::testing::{BAR, BAZ, FOO}; use test_case::test_case; + const PRIO: i32 = -100; + const REPO: &str = "owner1/repo1"; + + #[test_case("", &[""], &["a"]; "empty pattern")] + #[test_case( + "foo.bar+", + &["foo.bar+"], + &["fooXbar+", "foo.bar"]; + "regex metacharacters" + )] + #[test_case( + "foo%%_%bar", + &["foo-xbar", "foo-long-xbar"], + &["foobar", "foo-bar-baz"]; + "mixed consecutive wildcards" + )] + #[test_case( + r"literal\%percent", + &["literal%percent"], + &["literal-percent", "literalXXpercent"]; + "escaped percent" + )] + #[test_case( + r"literal\_underscore", + &["literal_underscore"], + &["literal-underscore", "literalXunderscore"]; + "escaped underscore" + )] + #[test_case( + r"literal\\backslash", + &[r"literal\backslash"], + &["literal/backslash", r"literal\\backslash"]; + "escaped backslash" + )] + #[test_case( + "trailing\\", + &["trailing\\"], + &["trailing", "trailing/"]; + "trailing backslash" + )] + #[test_case( + "cranelift-%", + &["cranelift-asdf", "cranelift-asdf-fb"], + &["other-xx" ]; + "prod example 1" + )] + #[test_case( + r"azure\_mgmt\_%", + &["azure_mgmt_123", "azure_mgmt_abc"], + &["azure-mgmt-123" ]; + "prod example 2" + )] + #[test_case("_", &["é", "🦀"], &["", "éé"]; "unicode wildcard")] + fn compile_like_pattern_handles_edge_cases( + pattern: &str, + should_match: &[&str], + should_not_match: &[&str], + ) -> Result<()> { + let regex = compile_like_pattern(pattern)?; + + for value in should_match { + assert!(regex.is_match(value), "{pattern:?} should match {value:?}"); + } + + for value in should_not_match { + assert!( + !regex.is_match(value), + "{pattern:?} should not match {value:?}" + ); + } + + Ok(()) + } + #[test_case( "docsrs-%", &["docsrs-database", "docsrs-", "docsrs-s3", "docsrs-webserver"], @@ -100,27 +333,22 @@ mod tests { should_match: &[&str], should_not_match: &[&str], ) -> Result<()> { - let test_metrics = TestMetrics::new(); - let db = TestDatabase::new(&Config::test_config()?, test_metrics.provider()).await?; + let env = TestEnv::new().await?; - const PRIO: i32 = -100; - - let mut conn = db.async_conn().await?; + let mut conn = env.db.async_conn().await?; set_crate_priority(&mut conn, pattern, PRIO).await?; + let priorities = PrioritiesCache::new(env.db.pool().clone(), 20); + for name in should_match { - assert_eq!( - get_crate_priority(&mut conn, &name.parse().unwrap()).await?, - PRIO - ); + let krate: KrateName = name.parse().unwrap(); + assert_eq!(priorities.get(&krate).await?, PRIO); } for name in should_not_match { - assert_eq!( - get_crate_priority(&mut conn, &name.parse().unwrap()).await?, - PRIORITY_DEFAULT - ); + let krate: KrateName = name.parse().unwrap(); + assert_eq!(priorities.get(&krate).await?, PRIORITY_DEFAULT); } Ok(()) @@ -128,42 +356,126 @@ mod tests { #[tokio::test(flavor = "multi_thread")] async fn remove_priority() -> Result<()> { - let test_metrics = TestMetrics::new(); - let db = TestDatabase::new(&Config::test_config()?, test_metrics.provider()).await?; + let env = TestEnv::new().await?; - let mut conn = db.async_conn().await?; + let mut conn = env.db.async_conn().await?; let pattern = "docsrs-%"; let krate = KrateName::from_static("docsrs-"); - const PRIO: i32 = -100; set_crate_priority(&mut conn, pattern, PRIO).await?; - assert_eq!(get_crate_priority(&mut conn, &krate).await?, PRIO); + let priorities = PrioritiesCache::new(env.db.pool().clone(), 20); + assert_eq!(priorities.get(&krate).await?, PRIO); assert_eq!(remove_crate_priority(&mut conn, pattern).await?, Some(PRIO)); - assert_eq!( - get_crate_priority(&mut conn, &krate).await?, - PRIORITY_DEFAULT - ); + priorities.reload().await?; + assert_eq!(priorities.get(&krate).await?, PRIORITY_DEFAULT); Ok(()) } #[tokio::test(flavor = "multi_thread")] async fn get_default_priority() -> Result<()> { - let test_metrics = TestMetrics::new(); - let db = TestDatabase::new(&Config::test_config()?, test_metrics.provider()).await?; + let env = TestEnv::new().await?; - let mut conn = db.async_conn().await?; + let priorities = PrioritiesCache::new(env.db.pool().clone(), 20); for name in &["docsrs", "rcc", "lasso", "hexponent", "rust4lyfe"] { let krate = KrateName::from_static(name); - assert_eq!( - get_crate_priority(&mut conn, &krate).await?, - PRIORITY_DEFAULT - ); + assert_eq!(priorities.get(&krate).await?, PRIORITY_DEFAULT); } Ok(()) } + + #[tokio::test(flavor = "multi_thread")] + async fn override_workspace_priority() -> Result<()> { + let env = TestEnv::new().await?; + + let mut conn = env.db.async_conn().await?; + + env.fake_release() + .await + .name(&FOO) + .github_stats_id( + FakeGithubStats::builder() + .repo(REPO) + .override_build_priority(-5) + .create(&mut conn) + .await?, + ) + .create() + .await?; + + assert_eq!( + get_crate_names_for_repository_build_priorities(&mut conn,).await?, + HashMap::from_iter([(FOO, -5)]) + ); + + let priorities = PrioritiesCache::new(env.db.pool().clone(), 20); + + // repo override is used + assert_eq!(priorities.get(&FOO).await?, -5); + + // no override, default prio + assert_eq!(priorities.get(&BAR).await?, PRIORITY_DEFAULT); + + // set pattern priority, should be used instead of the repo prio + set_crate_priority(&mut conn, FOO.as_ref(), PRIO).await?; + + priorities.reload().await?; + assert_eq!(priorities.get(&FOO).await?, PRIO); + + Ok(()) + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_auto_deprio_workspace() -> Result<()> { + let env = TestEnv::new().await?; + + let mut conn = env.db.async_conn().await?; + + let prioritized_repo_id = FakeGithubStats::builder() + .repo(REPO) + .create(&mut conn) + .await?; + + // two crates, one repo + for name in [FOO, BAR] { + env.fake_release() + .await + .name(&name) + .github_stats_id(prioritized_repo_id) + .create() + .await?; + } + + rewrite_repository_stats(&mut conn).await?; + + // validate our helper method, + assert!( + get_crate_names_from_big_workspaces(&mut conn, 3) + .await? + .is_empty() + ); + assert_eq!( + get_crate_names_from_big_workspaces(&mut conn, 2).await?, + HashSet::from_iter([FOO, BAR]) + ); + + let priorities = PrioritiesCache::new(env.db.pool().clone(), 2); + + // workspace size override is used + assert_eq!(priorities.get(&FOO).await?, PRIORITY_DEPRIORITIZED); + assert_eq!(priorities.get(&BAR).await?, PRIORITY_DEPRIORITIZED); + assert_eq!(priorities.get(&BAZ).await?, PRIORITY_DEFAULT); + + // set pattern priority, should be used instead of the workspace size prio + set_crate_priority(&mut conn, FOO.as_ref(), PRIO).await?; + + priorities.reload().await?; + assert_eq!(priorities.get(&FOO).await?, PRIO); + + Ok(()) + } } diff --git a/crates/lib/docs_rs_build_queue/src/queue/blocking.rs b/crates/lib/docs_rs_build_queue/src/queue/blocking.rs index 1554c0c285..a99313f5fc 100644 --- a/crates/lib/docs_rs_build_queue/src/queue/blocking.rs +++ b/crates/lib/docs_rs_build_queue/src/queue/blocking.rs @@ -11,8 +11,8 @@ use tracing::error; #[derive(Debug)] pub struct BuildQueue { - runtime: Handle, - inner: Arc, + pub(crate) runtime: Handle, + pub(crate) inner: Arc, } /// sync versions of async methods @@ -165,15 +165,10 @@ impl BuildQueue { #[cfg(test)] mod tests { - use crate::Config; - use super::*; + use crate::{Config, testing::test_env::BlockingTestEnv}; use chrono::Utc; - use docs_rs_config::AppConfig as _; - use docs_rs_database::{AsyncPoolClient, testing::TestDatabase}; - use docs_rs_opentelemetry::testing::TestMetrics; use docs_rs_types::testing::{KRATE, V1, V2}; - use docs_rs_utils::block_on_async_with_conn; use pretty_assertions::assert_eq; use std::time::Duration; @@ -181,84 +176,14 @@ mod tests { const BAR: KrateName = KrateName::from_static("bar"); const BAZ: KrateName = KrateName::from_static("baz"); - // when we start migrating / spitting the binaries, - // we probably will create amore powerfull & flexible - // test& app context. Then we could migrate this. - struct TestEnv { - db: TestDatabase, - queue: BuildQueue, - metrics: TestMetrics, - runtime: runtime::Runtime, - } - - impl TestEnv { - pub(crate) fn runtime(&self) -> &runtime::Runtime { - &self.runtime - } - - pub async fn async_conn(&self) -> Result { - self.db.async_conn().await - } - - fn queued_builds(&self) -> Result { - let collected_metrics = self.metrics.collected_metrics(); - - Ok(collected_metrics - .get_metric("build_queue", "docsrs.build_queue.queued_builds")? - .get_u64_counter() - .value()) - } - - fn failed_count(&self) -> u64 { - let collected_metrics = self.metrics.collected_metrics(); - - if let Ok(metric) = collected_metrics - .get_metric("build_queue", "docsrs.build_queue.failed_crates_count") - { - metric.get_u64_counter().value() - } else { - 0 - } - } - } - - fn test_queue(config: Config) -> Result { - let runtime = tokio::runtime::Builder::new_multi_thread() - .enable_all() - .build()?; - - let metrics = TestMetrics::new(); - let db = runtime.block_on(TestDatabase::new( - &docs_rs_database::Config::test_config()?, - metrics.provider(), - ))?; - - let async_queue = Arc::new(AsyncBuildQueue::new( - db.pool().clone(), - Arc::new(config), - metrics.provider(), - )); - - Ok(TestEnv { - db, - queue: BuildQueue { - runtime: runtime.handle().clone().into(), - inner: async_queue, - }, - metrics, - runtime, - }) - } - #[test] fn test_wait_between_build_attempts() -> Result<()> { - let env = test_queue(Config { + let env = BlockingTestEnv::new()?; + let queue = env.queue_with_config(Config { build_attempts: 99, delay_between_build_attempts: Duration::from_secs(1), ..Default::default() - })?; - - let queue = &env.queue; + }); queue.add_crate(&KRATE, &V1, 0)?; @@ -273,14 +198,16 @@ mod tests { unreachable!(); })?; - block_on_async_with_conn!(env, |mut conn| async { + env.block_on_async_with_conn(async |conn| { // fake the build-attempt timestamp so it's older - Ok(sqlx::query!( + sqlx::query!( "UPDATE queue SET last_attempt = $1", Utc::now() - chrono::Duration::try_seconds(60).unwrap() ) .execute(&mut *conn) - .await?) + .await?; + + Ok(()) })?; let mut handled = false; @@ -299,12 +226,12 @@ mod tests { #[test] fn test_add_and_process_crates() -> Result<()> { const MAX_ATTEMPTS: u16 = 3; - let env = test_queue(Config { + let env = BlockingTestEnv::new()?; + let queue = env.queue_with_config(Config { build_attempts: MAX_ATTEMPTS, delay_between_build_attempts: Duration::ZERO, ..Default::default() - })?; - let queue = &env.queue; + }); const LOW_PRIORITY: KrateName = KrateName::from_static("low-priority"); const HIGH_PRIORITY_FOO: KrateName = KrateName::from_static("high-priority-foo"); @@ -379,8 +306,8 @@ mod tests { #[test] fn test_pending_count() -> Result<()> { - let env = test_queue(Config::default())?; - let queue = env.queue; + let env = BlockingTestEnv::new()?; + let queue = env.queue(); assert_eq!(queue.pending_count()?, 0); queue.add_crate(&FOO, &V1, 0)?; assert_eq!(queue.pending_count()?, 1); @@ -398,8 +325,8 @@ mod tests { #[test] fn test_prioritized_count() -> Result<()> { - let env = test_queue(Config::default())?; - let queue = env.queue; + let env = BlockingTestEnv::new()?; + let queue = env.queue(); assert_eq!(queue.prioritized_count()?, 0); queue.add_crate(&FOO, &V1, 0)?; @@ -420,8 +347,8 @@ mod tests { #[test] fn test_count_by_priority() -> Result<()> { - let env = test_queue(Config::default())?; - let queue = env.queue; + let env = BlockingTestEnv::new()?; + let queue = env.queue(); assert!(queue.pending_count_by_priority()?.is_empty()); @@ -446,12 +373,12 @@ mod tests { fn test_failed_count_for_reattempts() -> Result<()> { const MAX_ATTEMPTS: u16 = 3; - let env = test_queue(Config { + let env = BlockingTestEnv::new()?; + let queue = env.queue_with_config(Config { build_attempts: MAX_ATTEMPTS, delay_between_build_attempts: Duration::ZERO, ..Default::default() - })?; - let queue = &env.queue; + }); assert_eq!(env.failed_count(), 0); queue.add_crate(&FOO, &V1, -100)?; @@ -483,12 +410,12 @@ mod tests { fn test_failed_count_after_error() -> Result<()> { const MAX_ATTEMPTS: u16 = 3; - let env = test_queue(Config { + let env = BlockingTestEnv::new()?; + let queue = env.queue_with_config(Config { build_attempts: MAX_ATTEMPTS, delay_between_build_attempts: Duration::ZERO, ..Default::default() - })?; - let queue = &env.queue; + }); assert_eq!(env.failed_count(), 0); queue.add_crate(&FOO, &V1, -100)?; @@ -515,8 +442,8 @@ mod tests { #[test] fn test_queued_crates() -> Result<()> { - let env = test_queue(Config::default())?; - let queue = env.queue; + let env = BlockingTestEnv::new()?; + let queue = env.queue(); let test_crates = [(BAR, 0), (FOO, -10), (BAZ, 10)]; for krate in &test_crates { @@ -537,8 +464,8 @@ mod tests { #[test] fn test_queue_lock() -> Result<()> { - let env = test_queue(Config::default())?; - let queue = env.queue; + let env = BlockingTestEnv::new()?; + let queue = env.queue(); // unlocked without config assert!(!queue.is_locked()?); @@ -554,8 +481,8 @@ mod tests { #[test] fn test_add_long_name() -> Result<()> { - let env = test_queue(Config::default())?; - let queue = env.queue; + let env = BlockingTestEnv::new()?; + let queue = env.queue(); let name: KrateName = "krate".repeat(100)[..64].parse().unwrap(); @@ -571,8 +498,8 @@ mod tests { #[test] fn test_add_long_version() -> Result<()> { - let env = test_queue(Config::default())?; - let queue = env.queue; + let env = BlockingTestEnv::new()?; + let queue = env.queue(); let long_version = Version::parse(&format!( "1.2.3-{}+{}", diff --git a/crates/lib/docs_rs_build_queue/src/queue/non_blocking.rs b/crates/lib/docs_rs_build_queue/src/queue/non_blocking.rs index ea85023bd2..ed512fb4da 100644 --- a/crates/lib/docs_rs_build_queue/src/queue/non_blocking.rs +++ b/crates/lib/docs_rs_build_queue/src/queue/non_blocking.rs @@ -1,6 +1,5 @@ use crate::{ - Config, PRIORITY_DEFAULT, PRIORITY_DEPRIORITIZED, PRIORITY_MANUAL_FROM_CRATES_IO, QueuedCrate, - metrics, + Config, PRIORITY_MANUAL_FROM_CRATES_IO, QueuedCrate, metrics, priority::PrioritiesCache, }; use anyhow::{Context as _, Result}; use docs_rs_database::{ @@ -8,31 +7,84 @@ use docs_rs_database::{ service_config::{Abnormality, ConfigName, get_config, set_config}, }; use docs_rs_opentelemetry::AnyMeterProvider; -use docs_rs_repository_stats::workspaces; use docs_rs_types::{KrateName, Version}; use docs_rs_uri::EscapedURI; use futures_util::TryStreamExt as _; -use std::{ - collections::{HashMap, HashSet}, - sync::Arc, -}; +use std::{collections::HashMap, sync::Arc}; #[derive(Debug)] pub struct AsyncBuildQueue { pub(super) config: Arc, pub(super) db: Pool, pub(super) queue_metrics: metrics::BuildQueueMetrics, + pub(super) priorities_cache: PrioritiesCache, } impl AsyncBuildQueue { pub fn new(db: Pool, config: Arc, otel_meter_provider: &AnyMeterProvider) -> Self { AsyncBuildQueue { + priorities_cache: PrioritiesCache::new(db.clone(), config.deprioritize_workspace_size), config, db, queue_metrics: metrics::BuildQueueMetrics::new(otel_meter_provider), } } + pub async fn find_priority(&self, name: &KrateName) -> Result { + self.priorities_cache.get(name).await + } + + /// Reevaluate priorities in build queue. + /// + /// Applies eventually updated priority configuration to already queued releases. + /// Non-default priorities are deliberately left unchanged. + pub async fn reevaluate_priorities(&self) -> Result<()> { + let mut conn = self.db.get_async().await?; + let names = sqlx::query_scalar!( + r#" + SELECT DISTINCT name AS "name: KrateName" + FROM queue + WHERE priority = $1 + "#, + crate::PRIORITY_DEFAULT, + ) + .fetch_all(&mut *conn) + .await?; + + let mut changed_names = Vec::new(); + let mut changed_priorities = Vec::new(); + + for name in names { + let priority = self.find_priority(&name).await?; + if priority != crate::PRIORITY_DEFAULT { + changed_names.push(name.to_string()); + changed_priorities.push(priority); + } + } + + if changed_names.is_empty() || changed_priorities.is_empty() { + return Ok(()); + } + + sqlx::query!( + r#" + UPDATE queue + SET priority = updates.priority + FROM UNNEST($1::text[], $2::int[]) AS updates(name, priority) + WHERE + queue.name = updates.name + AND queue.priority = $3 + "#, + &changed_names, + &changed_priorities, + crate::PRIORITY_DEFAULT, + ) + .execute(&mut *conn) + .await?; + + Ok(()) + } + pub async fn add_crate( &self, name: &KrateName, @@ -168,74 +220,6 @@ impl AsyncBuildQueue { Ok(()) } - /// Decreases the priority of all releases coming from bigger workspaces. - /// - /// In a separate method so we can do a bulk-select & update. - pub async fn deprioritize_workspaces(&self) -> Result<()> { - let mut conn = self.db.get_async().await?; - - let mut high_priority_names: HashSet<_> = sqlx::query!( - r#" - SELECT DISTINCT name AS "name: KrateName" - FROM queue - WHERE priority = $1"#, - PRIORITY_DEFAULT - ) - .fetch_all(&mut *conn) - .await? - .into_iter() - .map(|row| row.name) - .collect(); - - for (name, prio) in - workspaces::get_overriden_build_priorities(&mut conn, high_priority_names.iter()) - .await? - { - if prio != PRIORITY_DEFAULT { - sqlx::query!( - r#" - UPDATE queue - SET priority = $2 - WHERE name = $1 - "#, - name as _, - prio - ) - .execute(&mut *conn) - .await?; - } - - high_priority_names.remove(&name); - } - - let to_deprioritize: Vec<_> = - workspaces::get_crate_counts(&mut conn, high_priority_names.iter()) - .await? - .into_iter() - .filter_map(|(name, count)| { - (count > self.config.deprioritize_workspace_size.into()) - .then(|| name.to_string()) - }) - .collect(); - - sqlx::query!( - r#" - UPDATE queue - SET priority = $3 - WHERE - name = ANY($1) AND - priority = $2 - "#, - &to_deprioritize[..], - PRIORITY_DEFAULT, - PRIORITY_DEPRIORITIZED, - ) - .execute(&mut *conn) - .await?; - - Ok(()) - } - /// Decreases the priority of all releases currently present in the queue not matching the version passed to *at least* new_priority. pub async fn deprioritize_other_releases( &self, @@ -326,60 +310,24 @@ impl AsyncBuildQueue { #[cfg(test)] mod tests { + use crate::testing::test_env::TestEnv; + use super::*; use docs_rs_config::AppConfig as _; - use docs_rs_database::testing::TestDatabase; - use docs_rs_opentelemetry::testing::TestMetrics; - use docs_rs_storage::testing::TestStorage; - use docs_rs_test_fakes::{FakeGithubStats, FakeRelease}; + use docs_rs_repository_stats::workspaces::{ + rewrite_repository_stats, set_repository_build_priority, + }; + use docs_rs_test_fakes::FakeGithubStats; use docs_rs_types::testing::{BAR, BAZ, FOO, KRATE, V1, V2}; use pretty_assertions::assert_eq; const FAILED_KRATE: KrateName = KrateName::from_static("failed_crate"); - - // when we start migrating / spitting the binaries, - // we probably will create amore powerfull & flexible - // test& app context. Then we could migrate this. - struct TestEnv { - db: TestDatabase, - storage: TestStorage, - queue: AsyncBuildQueue, - } - - impl TestEnv { - async fn fake_release(&self) -> FakeRelease<'_> { - FakeRelease::new(self.db.pool().clone(), self.storage.storage().clone()) - } - } - - async fn test_queue() -> Result { - test_queue_with_config(Config::from_environment()?).await - } - - async fn test_queue_with_config(config: Config) -> Result { - let test_metrics = TestMetrics::new(); - let db = TestDatabase::new( - &docs_rs_database::Config::test_config()?, - test_metrics.provider(), - ) - .await?; - - let storage = TestStorage::from_config( - docs_rs_storage::Config::test_config()?.into(), - test_metrics.provider(), - ) - .await?; - - let queue = - AsyncBuildQueue::new(db.pool().clone(), Arc::new(config), test_metrics.provider()); - - Ok(TestEnv { db, storage, queue }) - } + const REPO: &str = "owner1/repo1"; #[tokio::test(flavor = "multi_thread")] async fn test_add_duplicate_doesnt_fail_last_priority_wins() -> Result<()> { - let env = test_queue().await?; - let queue = env.queue; + let env = TestEnv::new().await?; + let queue = env.queue(); queue.add_crate(&KRATE, &V1, 0).await?; queue.add_crate(&KRATE, &V1, 9).await?; @@ -393,8 +341,8 @@ mod tests { #[tokio::test(flavor = "multi_thread")] async fn test_add_duplicate_resets_attempts_and_priority() -> Result<()> { - let env = test_queue().await?; - let queue = env.queue; + let env = TestEnv::new().await?; + let queue = env.queue(); assert_eq!(queue.pending_count().await?, 0); @@ -430,10 +378,94 @@ mod tests { Ok(()) } + #[tokio::test(flavor = "multi_thread")] + async fn test_reevaluate_priorities_deprioritizes_queued_workspace() -> Result<()> { + let mut config = Config::from_environment()?; + config.deprioritize_workspace_size = 1; + let env = TestEnv::new().await?; + let queue = env.queue_with_config(config); + let mut conn = env.db.async_conn().await?; + + let repo_id = FakeGithubStats::builder() + .repo(REPO) + .create(&mut conn) + .await?; + + for name in [FOO, BAR] { + env.fake_release() + .await + .name(&name) + .version(V1) + .github_stats_id(repo_id) + .create() + .await?; + } + + for name in [FOO, BAR, BAZ] { + queue.add_crate(&name, &V1, crate::PRIORITY_DEFAULT).await?; + } + + rewrite_repository_stats(&mut conn).await?; + queue.priorities_cache.reload().await?; + queue.reevaluate_priorities().await?; + + assert_eq!( + queue + .queued_crates() + .await? + .into_iter() + .map(|queued| (queued.name, queued.priority)) + .collect::>(), + vec![ + (BAZ, crate::PRIORITY_DEFAULT), + (FOO, crate::PRIORITY_DEPRIORITIZED), + (BAR, crate::PRIORITY_DEPRIORITIZED), + ] + ); + + Ok(()) + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_reevaluate_priorities_applies_repository_override() -> Result<()> { + let env = TestEnv::new().await?; + let queue = env.queue(); + let mut conn = env.db.async_conn().await?; + + env.fake_release() + .await + .name(&FOO) + .version(V1) + .github_stats(REPO, 0, 0, 0) + .create() + .await?; + + queue.add_crate(&FOO, &V1, crate::PRIORITY_DEFAULT).await?; + queue + .add_crate(&BAR, &V1, crate::PRIORITY_MANUAL_FROM_CRATES_IO) + .await?; + + set_repository_build_priority(&mut conn, REPO, -10).await?; + queue.priorities_cache.reload().await?; + queue.reevaluate_priorities().await?; + + assert_eq!( + queue + .queued_crates() + .await? + .into_iter() + .map(|queued| (queued.name, queued.priority)) + .collect::>(), + vec![(FOO, -10), (BAR, crate::PRIORITY_MANUAL_FROM_CRATES_IO),] + ); + + Ok(()) + } + #[tokio::test(flavor = "multi_thread")] async fn test_has_build_queued() -> Result<()> { - let env = test_queue().await?; - let queue = env.queue; + let env = TestEnv::new().await?; + let queue = env.queue(); queue.add_crate(&KRATE, &V1, 0).await?; @@ -452,8 +484,8 @@ mod tests { #[tokio::test(flavor = "multi_thread")] async fn test_delete_version_from_queue() -> Result<()> { - let env = test_queue().await?; - let queue = env.queue; + let env = TestEnv::new().await?; + let queue = env.queue(); assert_eq!(queue.pending_count().await?, 0); @@ -478,8 +510,8 @@ mod tests { #[tokio::test(flavor = "multi_thread")] async fn test_delete_crate_from_queue() -> Result<()> { - let env = test_queue().await?; - let queue = env.queue; + let env = TestEnv::new().await?; + let queue = env.queue(); assert_eq!(queue.pending_count().await?, 0); @@ -494,139 +526,12 @@ mod tests { Ok(()) } - #[tokio::test(flavor = "multi_thread")] - async fn test_deprio_workspace() -> Result<()> { - let mut config = Config::from_environment()?; - // so any workspace with more than 1 crate is deprioritized - config.deprioritize_workspace_size = 1; - let env = test_queue_with_config(config).await?; - - let mut conn = env.db.async_conn().await?; - - // make FOO and BAR a workspace - { - let repo_id = FakeGithubStats { - repo: "owner1/repo1".into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; - - for name in [FOO, BAR] { - env.fake_release() - .await - .name(&name) - .version(V1) - .create() - .await?; - } - - sqlx::query!("UPDATE releases SET repository_id = $1", repo_id) - .execute(&mut *conn) - .await?; - } - - let queue = env.queue; - - // enqueue FOO and BAR with priority 0 - for krate in &[FOO, BAR, BAZ] { - queue.add_crate(krate, &V1, 0).await?; - } - - // unchanged priorities - assert_eq!( - queue - .queued_crates() - .await? - .into_iter() - .map(|q| (q.name, q.priority)) - .collect::>(), - vec![(FOO, 0), (BAR, 0), (BAZ, 0)] - ); - - workspaces::rewrite_repository_stats(&mut conn).await?; - - assert_eq!( - workspaces::get_crate_counts(&mut conn, [FOO, BAR, BAZ].iter()).await?, - HashMap::from_iter([(FOO, 2), (BAR, 2)]) - ); - - queue.deprioritize_workspaces().await?; - - // updated priorities & order - assert_eq!( - queue - .queued_crates() - .await? - .into_iter() - .map(|q| (q.name, q.priority)) - .collect::>(), - vec![(BAZ, 0), (FOO, 1), (BAR, 1)] - ); - - Ok(()) - } - - #[tokio::test(flavor = "multi_thread")] - async fn test_deprio_workspace_respects_repository_override_priority() -> Result<()> { - let mut config = Config::from_environment()?; - config.deprioritize_workspace_size = 1; - let env = test_queue_with_config(config).await?; - - let mut conn = env.db.async_conn().await?; - - let repo_id = FakeGithubStats { - repo: "owner1/repo1".into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; - - for name in [FOO, BAR] { - env.fake_release() - .await - .name(&name) - .version(V1) - .create() - .await?; - } - - sqlx::query!("UPDATE releases SET repository_id = $1", repo_id) - .execute(&mut *conn) - .await?; - workspaces::set_repository_build_priority(&mut conn, "owner1/repo1", -10).await?; - - let queue = env.queue; - for krate in &[FOO, BAR, BAZ] { - queue.add_crate(krate, &V1, PRIORITY_DEFAULT).await?; - } - - workspaces::rewrite_repository_stats(&mut conn).await?; - queue.deprioritize_workspaces().await?; - - assert_eq!( - queue - .queued_crates() - .await? - .into_iter() - .map(|q| (q.name, q.priority)) - .collect::>(), - vec![(FOO, -10), (BAR, -10), (BAZ, PRIORITY_DEFAULT)] - ); - - Ok(()) - } - #[tokio::test(flavor = "multi_thread")] async fn test_length_warning_threshold_boundary() -> Result<()> { let mut config = Config::from_environment()?; config.length_warning_threshold = 1; - let env = test_queue_with_config(config).await?; - let queue = env.queue; + let env = TestEnv::new().await?; + let queue = env.queue_with_config(config); queue.add_crate(&FOO, &V1, 0).await?; @@ -652,8 +557,8 @@ mod tests { async fn test_public_alert_ignores_manual_crates() -> Result<()> { let mut config = Config::from_environment()?; config.length_warning_threshold = 0; - let env = test_queue_with_config(config).await?; - let queue = env.queue; + let env = TestEnv::new().await?; + let queue = env.queue_with_config(config); queue .add_crate(&FOO, &V1, PRIORITY_MANUAL_FROM_CRATES_IO) diff --git a/crates/lib/docs_rs_build_queue/src/testing/mod.rs b/crates/lib/docs_rs_build_queue/src/testing/mod.rs new file mode 100644 index 0000000000..cb494ef96a --- /dev/null +++ b/crates/lib/docs_rs_build_queue/src/testing/mod.rs @@ -0,0 +1 @@ +pub(crate) mod test_env; diff --git a/crates/lib/docs_rs_build_queue/src/testing/test_env.rs b/crates/lib/docs_rs_build_queue/src/testing/test_env.rs new file mode 100644 index 0000000000..3a3592f532 --- /dev/null +++ b/crates/lib/docs_rs_build_queue/src/testing/test_env.rs @@ -0,0 +1,117 @@ +use crate::{AsyncBuildQueue, BuildQueue, Config}; +use anyhow::Result; +use docs_rs_config::AppConfig as _; +use docs_rs_database::testing::TestDatabase; +use docs_rs_opentelemetry::testing::TestMetrics; +use docs_rs_storage::testing::TestStorage; +use docs_rs_test_fakes::FakeRelease; +use std::sync::Arc; +use tokio::runtime; + +pub(crate) struct TestEnv { + metrics: TestMetrics, + pub(crate) db: TestDatabase, + pub(crate) storage: TestStorage, +} + +impl TestEnv { + pub(crate) async fn fake_release(&self) -> FakeRelease<'_> { + FakeRelease::new(self.db.pool().clone(), self.storage.storage().clone()) + } + + pub(crate) async fn new() -> Result { + let metrics = TestMetrics::new(); + let db = TestDatabase::new( + &docs_rs_database::Config::test_config()?, + metrics.provider(), + ) + .await?; + + let storage = TestStorage::from_config( + docs_rs_storage::Config::test_config()?.into(), + metrics.provider(), + ) + .await?; + + Ok(TestEnv { + metrics, + db, + storage, + }) + } + + pub(crate) fn queue(&self) -> AsyncBuildQueue { + self.queue_with_config(Config::default()) + } + + pub(crate) fn queue_with_config(&self, config: Config) -> AsyncBuildQueue { + AsyncBuildQueue::new( + self.db.pool().clone(), + Arc::new(config), + self.metrics.provider(), + ) + } +} + +pub(crate) struct BlockingTestEnv { + inner: TestEnv, + #[allow(dead_code)] // we need to keep the runtime alive while using the inner environment + runtime: runtime::Runtime, +} + +impl BlockingTestEnv { + pub(crate) fn new() -> Result { + let runtime = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .build()?; + + Ok(BlockingTestEnv { + inner: runtime.block_on(TestEnv::new())?, + runtime, + }) + } + + pub(crate) fn queue(&self) -> BuildQueue { + self.queue_with_config(Config::default()) + } + + pub(crate) fn queue_with_config(&self, config: Config) -> BuildQueue { + let async_queue = self.inner.queue_with_config(config); + + BuildQueue { + runtime: self.runtime.handle().clone().into(), + inner: async_queue.into(), + } + } + + pub(crate) fn queued_builds(&self) -> Result { + let collected_metrics = self.inner.metrics.collected_metrics(); + + Ok(collected_metrics + .get_metric("build_queue", "docsrs.build_queue.queued_builds")? + .get_u64_counter() + .value()) + } + + pub(crate) fn failed_count(&self) -> u64 { + let collected_metrics = self.inner.metrics.collected_metrics(); + + if let Ok(metric) = + collected_metrics.get_metric("build_queue", "docsrs.build_queue.failed_crates_count") + { + metric.get_u64_counter().value() + } else { + 0 + } + } + + pub(crate) fn block_on_async_with_conn( + &self, + f: impl AsyncFnOnce(&mut sqlx::PgConnection) -> Result, + ) -> Result { + self.runtime.block_on(async { + let mut conn = self.inner.db.async_conn().await?; + f(&mut conn).await + }) + } +} diff --git a/crates/lib/docs_rs_repository_stats/src/workspaces.rs b/crates/lib/docs_rs_repository_stats/src/workspaces.rs index f3485273b0..9278fb7d2d 100644 --- a/crates/lib/docs_rs_repository_stats/src/workspaces.rs +++ b/crates/lib/docs_rs_repository_stats/src/workspaces.rs @@ -1,7 +1,7 @@ use anyhow::{Result, anyhow, bail}; use docs_rs_types::KrateName; use sqlx::Acquire as _; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; /// update the crate-count on each repository. /// @@ -65,42 +65,37 @@ pub async fn rewrite_repository_stats(conn: &mut sqlx::PgConnection) -> Result<( Ok(()) } -/// get the crate-count for the related workspace for a crate. -pub async fn get_crate_counts( +/// get all crate names where the workspace crate count is greater than the given value +pub async fn get_crate_names_from_big_workspaces( conn: &mut sqlx::PgConnection, - names: impl Iterator, -) -> Result> { - let names: Vec<_> = names.map(|k: &KrateName| k.to_string()).collect(); - + min_workspace_size: u16, +) -> Result> { Ok(sqlx::query!( r#" SELECT - c.name as "name: KrateName", - repo.crate_count + c.name as "name: KrateName" FROM crates AS c INNER JOIN releases AS r ON c.latest_version_id = r.id INNER JOIN repositories AS repo ON r.repository_id = repo.id - WHERE c.name = ANY($1) + WHERE + repo.crate_count >= $1 "#, - &names[..], + min_workspace_size as i32, ) .fetch_all(&mut *conn) .await? .into_iter() - .map(|row| (row.name, row.crate_count)) + .map(|row| row.name) .collect()) } -/// get the override-priorities for the given crate names -pub async fn get_overriden_build_priorities( +/// get all crate names where the build priority is overridden on a workspace level +pub async fn get_crate_names_for_repository_build_priorities( conn: &mut sqlx::PgConnection, - names: impl Iterator, ) -> Result> { - let names: Vec<_> = names.map(|k: &KrateName| k.to_string()).collect(); - Ok(sqlx::query!( r#" SELECT @@ -113,10 +108,8 @@ pub async fn get_overriden_build_priorities( INNER JOIN repositories AS repo ON r.repository_id = repo.id WHERE - c.name = ANY($1) AND repo.override_build_priority IS NOT NULL "#, - &names[..], ) .fetch_all(&mut *conn) .await? @@ -326,57 +319,44 @@ mod tests { let env = test_env().await?; let mut conn = env.db.async_conn().await?; - let repo_id = FakeGithubStats { - repo: "owner/repo".into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; + let repo_id = FakeGithubStats::builder() + .repo(REPO) + .create(&mut conn) + .await?; for name in &[FOO, BAR] { - env.fake_release().await.name(name).create().await?; + env.fake_release() + .await + .name(name) + .github_stats_id(repo_id) + .create() + .await?; } - sqlx::query!("UPDATE releases SET repository_id = $1", repo_id) - .execute(&mut *conn) - .await?; - // the stats should be 0, because neither the full rewrite or the single-repo update was // called - assert_eq!( - fetch_stats(&mut conn).await?, - vec![("owner/repo".into(), 0)] - ); + assert_eq!(fetch_stats(&mut conn).await?, vec![(REPO.into(), 0)]); rewrite_repository_stats(&mut conn).await?; // after the full rewrite, the count is correct - assert_eq!( - fetch_stats(&mut conn).await?, - vec![("owner/repo".into(), 2)] - ); + assert_eq!(fetch_stats(&mut conn).await?, vec![(REPO.into(), 2)]); - env.fake_release().await.name(&BAZ).create().await?; - sqlx::query!("UPDATE releases SET repository_id = $1", repo_id) - .execute(&mut *conn) + env.fake_release() + .await + .name(&BAZ) + .github_stats_id(repo_id) + .create() .await?; // after adding a release, the count is still 2, // because neither the full rewrite or the single-repo update was called - assert_eq!( - fetch_stats(&mut conn).await?, - vec![("owner/repo".into(), 2)] - ); + assert_eq!(fetch_stats(&mut conn).await?, vec![(REPO.into(), 2)]); update_repository_stats(&mut conn, repo_id).await?; // here we expect the count to be 3, because we called the single-repo update - assert_eq!( - fetch_stats(&mut conn).await?, - vec![("owner/repo".into(), 3),] - ); + assert_eq!(fetch_stats(&mut conn).await?, vec![(REPO.into(), 3),]); Ok(()) } @@ -386,37 +366,28 @@ mod tests { let env = test_env().await?; let mut conn = env.db.async_conn().await?; - let prioritized_repo_id = FakeGithubStats { - repo: REPO.into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; - - let default_repo_id = FakeGithubStats { - repo: REPO2.into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; - - env.fake_release().await.name(&FOO).create().await?; - env.fake_release().await.name(&BAR).create().await?; - - sqlx::query!("UPDATE releases SET repository_id = $1 WHERE crate_id = (SELECT id FROM crates WHERE name = $2)", prioritized_repo_id, FOO as _) - .execute(&mut *conn) + env.fake_release() + .await + .name(&FOO) + .github_stats_id( + FakeGithubStats::builder() + .repo(REPO) + .override_build_priority(-5) + .create(&mut conn) + .await?, + ) + .create() .await?; - sqlx::query!("UPDATE releases SET repository_id = $1 WHERE crate_id = (SELECT id FROM crates WHERE name = $2)", default_repo_id, BAR as _) - .execute(&mut *conn) + + env.fake_release() + .await + .name(&BAR) + .github_stats(REPO2, 0, 0, 0) + .create() .await?; - set_repository_build_priority(&mut conn, REPO, -5).await?; assert_eq!( - get_overriden_build_priorities(&mut conn, [FOO, BAR, BAZ].iter()).await?, + get_crate_names_for_repository_build_priorities(&mut conn).await?, HashMap::from_iter([(FOO, -5)]) ); @@ -428,14 +399,10 @@ mod tests { let env = test_env().await?; let mut conn = env.db.async_conn().await?; - FakeGithubStats { - repo: REPO.into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; + FakeGithubStats::builder() + .repo(REPO) + .create(&mut conn) + .await?; assert_eq!(get_repository_build_priority(&mut conn, REPO).await?, None); @@ -500,25 +467,16 @@ mod tests { let env = test_env().await?; let mut conn = env.db.async_conn().await?; - FakeGithubStats { - repo: REPO.into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; - - FakeGithubStats { - repo: REPO2.into(), - stars: 0, - forks: 0, - issues: 0, - } - .create(&mut conn) - .await?; + FakeGithubStats::builder() + .repo(REPO) + .override_build_priority(-5) + .create(&mut conn) + .await?; - set_repository_build_priority(&mut conn, REPO, -5).await?; + FakeGithubStats::builder() + .repo(REPO2) + .create(&mut conn) + .await?; assert_eq!( list_repository_build_priorities(&mut conn).await?, diff --git a/crates/lib/docs_rs_test_fakes/Cargo.toml b/crates/lib/docs_rs_test_fakes/Cargo.toml index 5915e374a5..055ae01469 100644 --- a/crates/lib/docs_rs_test_fakes/Cargo.toml +++ b/crates/lib/docs_rs_test_fakes/Cargo.toml @@ -8,6 +8,7 @@ edition = "2024" [dependencies] anyhow = { workspace = true } base64 = { workspace = true } +bon = { workspace = true } chrono = { workspace = true } docs_rs_cargo_metadata = { path = "../docs_rs_cargo_metadata", features = ["testing"] } docs_rs_database = { path = "../docs_rs_database" } diff --git a/crates/lib/docs_rs_test_fakes/src/github_stats.rs b/crates/lib/docs_rs_test_fakes/src/github_stats.rs new file mode 100644 index 0000000000..927f485e3f --- /dev/null +++ b/crates/lib/docs_rs_test_fakes/src/github_stats.rs @@ -0,0 +1,78 @@ +use anyhow::Result; +use base64::{Engine, engine::general_purpose::STANDARD as b64}; + +#[derive(bon::Builder)] +#[builder(on(_, into))] +pub struct FakeGithubStats { + repo: String, + + #[builder(default)] + stars: i32, + + #[builder(default)] + forks: i32, + + #[builder(default)] + issues: i32, + + override_build_priority: Option, +} + +impl FakeGithubStats { + pub async fn create(&self, conn: &mut sqlx::PgConnection) -> Result { + let existing_count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM repositories") + .fetch_one(&mut *conn) + .await? + .unwrap(); + let host_id = b64.encode(format!("FAKE ID {existing_count}")); + + let id = sqlx::query_scalar!( + "INSERT INTO repositories ( + host, + host_id, + name, + description, + last_commit, + stars, + forks, + issues, + updated_at, + override_build_priority + ) + VALUES ( + 'github.com', + $1, + $2, + 'Fake description!', + NOW(), + $3, + $4, + $5, + NOW(), + $6 + ) + RETURNING id", + host_id, + self.repo, + self.stars, + self.forks, + self.issues, + self.override_build_priority, + ) + .fetch_one(&mut *conn) + .await?; + + Ok(id) + } +} + +use fake_github_stats_builder::{IsComplete, State}; + +impl FakeGithubStatsBuilder { + pub async fn create(self, conn: &mut sqlx::PgConnection) -> Result + where + S: IsComplete, + { + self.build().create(conn).await + } +} diff --git a/crates/lib/docs_rs_test_fakes/src/legacy.rs b/crates/lib/docs_rs_test_fakes/src/legacy.rs index be9460bfe3..5fe184cd57 100644 --- a/crates/lib/docs_rs_test_fakes/src/legacy.rs +++ b/crates/lib/docs_rs_test_fakes/src/legacy.rs @@ -1,5 +1,5 @@ +use crate::FakeGithubStats; use anyhow::{Context as _, Result, bail}; -use base64::{Engine, engine::general_purpose::STANDARD as b64}; use chrono::{DateTime, Utc}; use docs_rs_cargo_metadata::{Dependency, MetadataPackage, Target}; use docs_rs_database::{ @@ -86,6 +86,7 @@ pub struct FakeRelease<'a> { /// This stores the content, while `package.readme` stores the filename readme: Option<&'a str>, github_stats: Option, + github_stats_id: Option, doc_coverage: Option, no_cargo_toml: bool, } @@ -155,6 +156,7 @@ impl<'a> FakeRelease<'a> { has_examples: false, readme: None, github_stats: None, + github_stats_id: None, doc_coverage: None, archive_storage: false, no_cargo_toml: false, @@ -339,12 +341,19 @@ impl<'a> FakeRelease<'a> { forks: i32, issues: i32, ) -> Self { - self.github_stats = Some(FakeGithubStats { - repo: repo.into(), - stars, - forks, - issues, - }); + self.github_stats = Some( + FakeGithubStats::builder() + .repo(repo) + .stars(stars) + .forks(forks) + .issues(issues) + .build(), + ); + self + } + + pub fn github_stats_id(mut self, id: i32) -> Self { + self.github_stats_id = Some(id); self } @@ -519,9 +528,13 @@ impl<'a> FakeRelease<'a> { let mut async_conn = pool.get_async().await?; - let repository = match self.github_stats { - Some(stats) => Some(stats.create(&mut async_conn).await?), - None => None, + let repository = match (self.github_stats, self.github_stats_id) { + (Some(_), Some(_)) => { + bail!("can't have both given github stats and an external github stats id") + } + (Some(stats), None) => Some(stats.create(&mut async_conn).await?), + (None, Some(id)) => Some(id), + (None, None) => None, }; let crate_tmp = create_temp_dir(); @@ -611,32 +624,6 @@ impl<'a> FakeRelease<'a> { } } -pub struct FakeGithubStats { - pub repo: String, - pub stars: i32, - pub forks: i32, - pub issues: i32, -} - -impl FakeGithubStats { - pub async fn create(&self, conn: &mut sqlx::PgConnection) -> Result { - let existing_count: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM repositories") - .fetch_one(&mut *conn) - .await? - .unwrap(); - let host_id = b64.encode(format!("FAKE ID {existing_count}")); - - let id = sqlx::query_scalar!( - "INSERT INTO repositories (host, host_id, name, description, last_commit, stars, forks, issues, updated_at) - VALUES ('github.com', $1, $2, 'Fake description!', NOW(), $3, $4, $5, NOW()) - RETURNING id", - host_id, self.repo, self.stars, self.forks, self.issues, - ).fetch_one(&mut *conn).await?; - - Ok(id) - } -} - impl FakeBuild { pub fn rustc_version(self, rustc_version: impl Into) -> Self { Self { diff --git a/crates/lib/docs_rs_test_fakes/src/lib.rs b/crates/lib/docs_rs_test_fakes/src/lib.rs index ed41225690..b52b61eb4a 100644 --- a/crates/lib/docs_rs_test_fakes/src/lib.rs +++ b/crates/lib/docs_rs_test_fakes/src/lib.rs @@ -1,4 +1,6 @@ +mod github_stats; mod legacy; pub use docs_rs_registry_api::{CrateOwner, OwnerKind}; -pub use legacy::{FakeBuild, FakeGithubStats, FakeRelease, fake_release_that_failed_before_build}; +pub use github_stats::FakeGithubStats; +pub use legacy::{FakeBuild, FakeRelease, fake_release_that_failed_before_build};