diff --git a/apps/frontend/src/pages/moderation/technical-review/rules.vue b/apps/frontend/src/pages/moderation/technical-review/rules.vue index 2a5b86f2ec..d0447852c7 100644 --- a/apps/frontend/src/pages/moderation/technical-review/rules.vue +++ b/apps/frontend/src/pages/moderation/technical-review/rules.vue @@ -133,6 +133,14 @@ proceed-label="Delete rule" @proceed="deleteRule" /> +
@@ -148,14 +156,45 @@
- - - +
+ + + + + + +
+
+
+
+

Scanning Delphi rule effects

+

+ {{ scanProgress.scanned.toLocaleString() }} of + {{ scanProgress.total.toLocaleString() }} details scanned ยท + {{ scanProgress.effects.toLocaleString() }} effects +

+
+ + {{ scanProgress.phase }} revision {{ scanProgress.revision }} + +
+ +
+
CEL contract and input
@@ -164,6 +203,9 @@ { "severity": "low", "hidden": false }. Severity can be low, medium, high, or severe.

+

+ Rules run in the order shown below. The first rule that returns a non-null effect wins. +

The input object contains schema_version, trace (key, issue_type, severity, @@ -195,17 +237,17 @@

{{ rule.name }}

-

Last applied in revision {{ rule.revision }}

+

Revision {{ rule.revision }}

- - @@ -221,12 +263,13 @@ diff --git a/apps/labrinth/src/routes/internal/mod.rs b/apps/labrinth/src/routes/internal/mod.rs index 46cd42c713..d54c86c742 100644 --- a/apps/labrinth/src/routes/internal/mod.rs +++ b/apps/labrinth/src/routes/internal/mod.rs @@ -120,6 +120,7 @@ pub fn config(cfg: &mut web::ServiceConfig) { moderation::tech_review::rules::create_rule, moderation::tech_review::rules::update_rule, moderation::tech_review::rules::delete_rule, + moderation::tech_review::rules_scan::scan_rules, moderation::tech_review::get_project_report, moderation::tech_review::submit_report, moderation::tech_review::update_issue_details, diff --git a/apps/labrinth/src/routes/internal/moderation/tech_review.rs b/apps/labrinth/src/routes/internal/moderation/tech_review.rs index fedf2288f2..92b230cd79 100644 --- a/apps/labrinth/src/routes/internal/moderation/tech_review.rs +++ b/apps/labrinth/src/routes/internal/moderation/tech_review.rs @@ -45,11 +45,13 @@ use eyre::eyre; pub mod global; pub mod rules; +pub mod rules_scan; pub fn config(cfg: &mut actix_web::web::ServiceConfig) { cfg.service(search_projects) .configure(global::config) .configure(rules::config) + .configure(rules_scan::config) .service(get_project_report) .service(get_report) .service(get_issue) diff --git a/apps/labrinth/src/routes/internal/moderation/tech_review/rules.rs b/apps/labrinth/src/routes/internal/moderation/tech_review/rules.rs index 415956caa2..9a0a55d23c 100644 --- a/apps/labrinth/src/routes/internal/moderation/tech_review/rules.rs +++ b/apps/labrinth/src/routes/internal/moderation/tech_review/rules.rs @@ -191,32 +191,12 @@ pub async fn test_rule( for (index, trace) in request.traces.iter().enumerate() { let input = test_rule_input(trace); - let mut context = cel::Context::default(); - context.add_variable("input", input).map_err(|error| { - ApiError::Request(eyre!( - "failed to build input for test trace {index}: {error}" - )) - })?; - - let value = program.execute(&context).map_err(|error| { - ApiError::Request(eyre!( - "failed to evaluate test trace {index}: {error}" - )) - })?; - let value = value.json().map_err(|error| { - ApiError::Request(eyre!( - "failed to decode result for test trace {index}: {error}" - )) - })?; - - let effect = match value { - serde_json::Value::Null => None, - value => Some(serde_json::from_value(value).map_err(|error| { + let effect = super::rules_scan::evaluate_rule(&program, input) + .map_err(|error| { ApiError::Request(eyre!( - "invalid effect for test trace {index}: {error}" + "failed to evaluate test trace {index}: {error}" )) - })?), - }; + })?; effects.push(effect); } @@ -280,7 +260,7 @@ pub async fn get_rules( updated_by FROM delphi_rules WHERE NOT delete_on_next_revision - ORDER BY name, id + ORDER BY id "#, ) .fetch_all(&***ro_pool) @@ -336,10 +316,17 @@ pub async fn create_rule( INSERT INTO delphi_rules ( name, rule, + revision, created_by, updated_by ) - VALUES ($1, $2, $3, $3) + VALUES ( + $1, + $2, + (SELECT revision + 1 FROM delphi_rule_revisions LIMIT 1), + $3, + $3 + ) RETURNING id, name, @@ -405,6 +392,9 @@ pub async fn update_rule( SET name = $2, rule = $3, + revision = ( + SELECT revision + 1 FROM delphi_rule_revisions LIMIT 1 + ), updated_at = CURRENT_TIMESTAMP, updated_by = $4 WHERE id = $1 AND NOT delete_on_next_revision @@ -470,6 +460,9 @@ pub async fn delete_rule( UPDATE delphi_rules SET delete_on_next_revision = TRUE, + revision = ( + SELECT revision + 1 FROM delphi_rule_revisions LIMIT 1 + ), updated_at = CURRENT_TIMESTAMP, updated_by = $2 WHERE id = $1 AND NOT delete_on_next_revision diff --git a/apps/labrinth/src/routes/internal/moderation/tech_review/rules_scan.rs b/apps/labrinth/src/routes/internal/moderation/tech_review/rules_scan.rs new file mode 100644 index 0000000000..4ee7f4c595 --- /dev/null +++ b/apps/labrinth/src/routes/internal/moderation/tech_review/rules_scan.rs @@ -0,0 +1,493 @@ +use std::collections::{BTreeMap, HashMap}; + +use actix_web::{HttpRequest, HttpResponse, post, web}; +use ariadne::ids::base62_impl::to_base62; +use bytes::Bytes; +use eyre::{Context as _, Result, eyre}; +use futures_util::{StreamExt, TryStreamExt}; +use serde::Serialize; +use sqlx::types::Json; +use tokio::sync::mpsc; +use tokio_stream::wrappers::UnboundedReceiverStream; + +use super::rules::DelphiRuleEffect; +use crate::{ + auth::check_is_moderator_from_headers, + database::{ + PgPool, PgTransaction, models::delphi_report_item::DelphiSeverity, + redis::RedisPool, + }, + models::pats::Scopes, + queue::session::AuthQueue, + routes::ApiError, +}; + +const RULE_SCAN_LOCK_ID: i64 = 0x6465_6c70_6869_7275; +const PROGRESS_INTERVAL: usize = 50; + +pub fn config(cfg: &mut actix_web::web::ServiceConfig) { + cfg.service(scan_rules); +} + +#[derive(Serialize)] +struct RuleScanEvent<'a> { + phase: &'a str, + revision: i64, + scanned: usize, + total: usize, + effects: usize, +} + +#[derive(Serialize)] +struct RuleScanErrorEvent<'a> { + message: &'a str, +} + +#[derive(Serialize)] +struct RuleInput { + schema_version: u32, + trace: RuleTrace, + scan: RuleScan, + artifact: RuleArtifact, + scope: RuleScope, +} + +#[derive(Serialize)] +struct RuleTrace { + key: String, + issue_type: String, + severity: DelphiSeverity, + jar: Option, + file_path: String, + data: HashMap, +} + +#[derive(Serialize)] +struct RuleScan { + delphi_version: i32, +} + +#[derive(Serialize)] +struct RuleArtifact { + size: Option, + hashes: BTreeMap, +} + +#[derive(Serialize)] +struct RuleScope { + project_id: Option, + version_id: Option, + file_id: Option, +} + +struct CompiledRule { + id: i64, + program: cel::Program, +} + +struct MaterializedEffect { + detail_id: i64, + rule_id: i64, + effect: DelphiRuleEffect, +} + +struct ScanSummary { + revision: i64, + scanned: usize, + total: usize, + effects: usize, +} + +/// Re-evaluate every Delphi issue detail and atomically publish a new rule revision. +#[utoipa::path( + context_path = "/moderation/tech-review", + tag = "moderation", + security(("bearer_auth" = [])), + responses((status = OK), (status = CONFLICT)) +)] +#[post("/rules/scan")] +pub async fn scan_rules( + req: HttpRequest, + pool: web::Data, + redis: web::Data, + session_queue: web::Data, +) -> Result { + check_is_moderator_from_headers( + &req, + &**pool, + &redis, + &session_queue, + Scopes::PROJECT_WRITE, + ) + .await?; + + let mut transaction = crate::util::error::Context::wrap_internal_err( + pool.begin().await, + "failed to begin delphi rule scan", + )?; + + sqlx::query!("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ") + .execute(&mut transaction) + .await + .map_err(|error| { + ApiError::Internal( + eyre!(error) + .wrap_err("failed to set delphi rule scan isolation"), + ) + })?; + + let acquired = sqlx::query_scalar!( + "SELECT pg_try_advisory_xact_lock($1)", + RULE_SCAN_LOCK_ID, + ) + .fetch_one(&mut transaction) + .await + .map_err(|error| { + ApiError::Internal( + eyre!(error).wrap_err("failed to acquire delphi rule scan lock"), + ) + })? + .unwrap_or(false); + + if !acquired { + return Err(ApiError::Conflict( + "a delphi rule scan is already running".to_string(), + )); + } + + let (sender, receiver) = mpsc::unbounded_channel(); + actix_web::rt::spawn(async move { + match run_scan(transaction, &sender).await { + Ok(summary) => { + send_event( + &sender, + "complete", + &RuleScanEvent { + phase: "complete", + revision: summary.revision, + scanned: summary.scanned, + total: summary.total, + effects: summary.effects, + }, + ); + } + Err(error) => { + tracing::error!(error = ?error, "delphi rule scan failed"); + send_event( + &sender, + "failed", + &RuleScanErrorEvent { + message: &error.to_string(), + }, + ); + } + } + }); + + let stream = + UnboundedReceiverStream::new(receiver).map(Ok::<_, std::io::Error>); + + Ok(HttpResponse::Ok() + .insert_header(("Content-Type", "text/event-stream")) + .insert_header(("Cache-Control", "no-cache")) + .insert_header(("X-Accel-Buffering", "no")) + .streaming(stream)) +} + +async fn run_scan( + mut transaction: PgTransaction<'static>, + sender: &mpsc::UnboundedSender, +) -> Result { + sqlx::query!("LOCK TABLE delphi_rules IN SHARE MODE") + .execute(&mut transaction) + .await + .wrap_err("failed to lock delphi rules")?; + sqlx::query!("LOCK TABLE delphi_report_issue_details IN SHARE MODE") + .execute(&mut transaction) + .await + .wrap_err("failed to lock delphi issue details")?; + + let current_revision = sqlx::query_scalar!( + "SELECT revision FROM delphi_rule_revisions LIMIT 1 FOR UPDATE", + ) + .fetch_one(&mut transaction) + .await + .wrap_err("failed to fetch the current delphi rule revision")?; + let revision = current_revision + .checked_add(1) + .ok_or_else(|| eyre!("delphi rule revision overflowed"))?; + + let rules = sqlx::query!( + r#" + SELECT id, rule + FROM delphi_rules + WHERE NOT delete_on_next_revision + ORDER BY id + "#, + ) + .fetch_all(&mut transaction) + .await + .wrap_err("failed to fetch delphi rules")? + .into_iter() + .map(|rule| { + let program = cel::Program::compile(&rule.rule).map_err(|error| { + eyre!("failed to compile delphi rule {}: {error}", rule.id) + })?; + Ok(CompiledRule { + id: rule.id, + program, + }) + }) + .collect::>>()?; + + let total = sqlx::query_scalar!( + "SELECT COUNT(*) AS \"count!\" FROM delphi_report_issue_details", + ) + .fetch_one(&mut transaction) + .await + .wrap_err("failed to count delphi issue details")? as usize; + + let mut details = sqlx::query!( + r#" + SELECT + detail.id, + detail.key, + issue.issue_type, + detail.severity AS "severity: DelphiSeverity", + detail.jar, + detail.file_path, + detail.data AS "data: Json>", + report.delphi_version, + file.size AS "size?", + file.id AS "file_id?", + version.id AS "version_id?", + version.mod_id AS "project_id?", + COALESCE(file_hashes.hashes, '{}'::jsonb) + AS "hashes!: Json>" + FROM delphi_report_issue_details detail + INNER JOIN delphi_report_issues issue ON issue.id = detail.issue_id + INNER JOIN delphi_reports report ON report.id = issue.report_id + LEFT JOIN files file ON file.id = report.file_id + LEFT JOIN versions version ON version.id = file.version_id + LEFT JOIN ( + SELECT + file_id, + jsonb_object_agg(algorithm, encode(hash, 'hex')) AS hashes + FROM hashes + GROUP BY file_id + ) file_hashes ON file_hashes.file_id = file.id + ORDER BY detail.id + "#, + ) + .fetch(&mut transaction); + + let mut effects = Vec::new(); + let mut scanned = 0; + send_progress(sender, "scanning", revision, 0, total, 0); + + while let Some(detail) = details + .try_next() + .await + .wrap_err("failed to fetch a delphi issue detail")? + { + let detail_id = detail.id; + let input = RuleInput { + schema_version: 1, + trace: RuleTrace { + key: detail.key, + issue_type: detail.issue_type, + severity: detail.severity, + jar: detail.jar, + file_path: detail.file_path, + data: detail.data.0, + }, + scan: RuleScan { + delphi_version: detail.delphi_version, + }, + artifact: RuleArtifact { + size: detail.size, + hashes: detail.hashes.0, + }, + scope: RuleScope { + project_id: detail.project_id.map(to_public_id), + version_id: detail.version_id.map(to_public_id), + file_id: detail.file_id.map(to_public_id), + }, + }; + + for rule in &rules { + let effect = evaluate_rule(&rule.program, &input).wrap_err_with(|| { + format!( + "failed to evaluate delphi rule {} for detail {detail_id}", + rule.id + ) + })?; + if let Some(effect) = effect { + effects.push(MaterializedEffect { + detail_id, + rule_id: rule.id, + effect, + }); + break; + } + } + + scanned += 1; + if scanned % PROGRESS_INTERVAL == 0 || scanned == total { + send_progress( + sender, + "scanning", + revision, + scanned, + total, + effects.len(), + ); + tokio::task::yield_now().await; + } + } + drop(details); + + send_progress(sender, "publishing", revision, total, total, effects.len()); + + let detail_ids = effects + .iter() + .map(|effect| effect.detail_id) + .collect::>(); + let rule_ids = effects + .iter() + .map(|effect| effect.rule_id) + .collect::>(); + let severities = effects + .iter() + .map(|effect| effect.effect.severity) + .collect::>(); + let hidden = effects + .iter() + .map(|effect| effect.effect.hidden) + .collect::>(); + + if !effects.is_empty() { + sqlx::query!( + r#" + INSERT INTO delphi_rule_effects ( + revision, + detail_id, + rule_id, + severity, + hidden + ) + SELECT $1, effect.* + FROM UNNEST( + $2::BIGINT[], + $3::BIGINT[], + $4::delphi_severity[], + $5::BOOLEAN[] + ) AS effect(detail_id, rule_id, severity, hidden) + "#, + revision, + &detail_ids, + &rule_ids, + &severities as &[Option], + &hidden, + ) + .execute(&mut transaction) + .await + .wrap_err("failed to insert delphi rule effects")?; + } + + sqlx::query!( + "DELETE FROM delphi_rule_effects WHERE revision <> $1", + revision, + ) + .execute(&mut transaction) + .await + .wrap_err("failed to delete old delphi rule effects")?; + sqlx::query!("DELETE FROM delphi_rules WHERE delete_on_next_revision") + .execute(&mut transaction) + .await + .wrap_err("failed to delete retired delphi rules")?; + sqlx::query!( + "UPDATE delphi_rules SET revision = $1 WHERE NOT delete_on_next_revision", + revision, + ) + .execute(&mut transaction) + .await + .wrap_err("failed to update delphi rule revisions")?; + sqlx::query!("UPDATE delphi_rule_revisions SET revision = $1", revision) + .execute(&mut transaction) + .await + .wrap_err("failed to publish the delphi rule revision")?; + + transaction + .commit() + .await + .wrap_err("failed to commit the delphi rule scan")?; + + Ok(ScanSummary { + revision, + scanned: total, + total, + effects: effects.len(), + }) +} + +pub(super) fn evaluate_rule( + program: &cel::Program, + input: impl Serialize, +) -> Result> { + let mut context = cel::Context::default(); + context + .add_variable("input", input) + .wrap_err("failed to build cel input")?; + + let value = program + .execute(&context) + .wrap_err("failed to execute cel expression")?; + let value = value.json().map_err(|error| { + eyre!("failed to convert cel result to json: {error}") + })?; + + match value { + serde_json::Value::Null => Ok(None), + value => serde_json::from_value(value) + .map(Some) + .wrap_err("cel expression returned an invalid rule effect"), + } +} + +fn to_public_id(id: i64) -> String { + to_base62(id as u64) +} + +fn send_progress( + sender: &mpsc::UnboundedSender, + phase: &'static str, + revision: i64, + scanned: usize, + total: usize, + effects: usize, +) { + send_event( + sender, + "progress", + &RuleScanEvent { + phase, + revision, + scanned, + total, + effects, + }, + ); +} + +fn send_event( + sender: &mpsc::UnboundedSender, + event: &str, + data: &impl Serialize, +) { + let Ok(data) = serde_json::to_string(data) else { + return; + }; + let _ = + sender.send(Bytes::from(format!("event: {event}\ndata: {data}\n\n"))); +} diff --git a/packages/api-client/src/modules/labrinth/tech-review/internal.ts b/packages/api-client/src/modules/labrinth/tech-review/internal.ts index c5d86fc713..c66520baa8 100644 --- a/packages/api-client/src/modules/labrinth/tech-review/internal.ts +++ b/packages/api-client/src/modules/labrinth/tech-review/internal.ts @@ -68,6 +68,15 @@ export class LabrinthTechReviewInternalModule extends AbstractModule { }) } + public async scanRules(signal?: AbortSignal): Promise> { + return this.client.stream('/moderation/tech-review/rules/scan', { + api: 'labrinth', + version: 'internal', + method: 'POST', + signal, + }) + } + /** * Search for projects awaiting technical review. * diff --git a/packages/api-client/src/modules/labrinth/types.ts b/packages/api-client/src/modules/labrinth/types.ts index 51f664faa3..9d97f1ec5a 100644 --- a/packages/api-client/src/modules/labrinth/types.ts +++ b/packages/api-client/src/modules/labrinth/types.ts @@ -2265,6 +2265,20 @@ export namespace Labrinth { effects: Array } + export type DelphiRuleScanPhase = 'scanning' | 'publishing' | 'complete' + + export type DelphiRuleScanEvent = { + phase: DelphiRuleScanPhase + revision: number + scanned: number + total: number + effects: number + } + + export type DelphiRuleScanErrorEvent = { + message: string + } + export type SearchProjectsRequest = { limit?: number page?: number