@@ -112,9 +184,7 @@
{{ rule.name }}
-
- Priority {{ rule.priority }} · revision {{ rule.revision_id }}
-
+
Last applied in revision {{ rule.revision }}
@@ -133,7 +203,7 @@
{{ rule.expression }}
+ >
{{ rule.rule }}
@@ -141,7 +211,14 @@
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