From 64b10310b4258ecfe5696c5169cf0fd2f2b3ae39 Mon Sep 17 00:00:00 2001 From: Trezy Date: Thu, 7 May 2026 11:43:48 -0500 Subject: [PATCH] feat: add collection column to dead_letter_scripts so we can ditch second-level filtering Signed-off-by: Trezy Signed-off-by: Trezy --- ...2000000_dead_letter_scripts_collection.sql | 4 + ...2000000_dead_letter_scripts_collection.sql | 4 + src/admin/dead_letters.rs | 73 +++++++++---------- src/lua/scripts.rs | 13 +++- 4 files changed, 54 insertions(+), 40 deletions(-) create mode 100644 migrations/postgres/20260502000000_dead_letter_scripts_collection.sql create mode 100644 migrations/sqlite/20260502000000_dead_letter_scripts_collection.sql diff --git a/migrations/postgres/20260502000000_dead_letter_scripts_collection.sql b/migrations/postgres/20260502000000_dead_letter_scripts_collection.sql new file mode 100644 index 0000000..fae7c05 --- /dev/null +++ b/migrations/postgres/20260502000000_dead_letter_scripts_collection.sql @@ -0,0 +1,4 @@ +ALTER TABLE dead_letter_scripts ADD COLUMN collection TEXT; + +CREATE INDEX idx_dead_letter_scripts_collection + ON dead_letter_scripts (collection, resolved_at); diff --git a/migrations/sqlite/20260502000000_dead_letter_scripts_collection.sql b/migrations/sqlite/20260502000000_dead_letter_scripts_collection.sql new file mode 100644 index 0000000..fae7c05 --- /dev/null +++ b/migrations/sqlite/20260502000000_dead_letter_scripts_collection.sql @@ -0,0 +1,4 @@ +ALTER TABLE dead_letter_scripts ADD COLUMN collection TEXT; + +CREATE INDEX idx_dead_letter_scripts_collection + ON dead_letter_scripts (collection, resolved_at); diff --git a/src/admin/dead_letters.rs b/src/admin/dead_letters.rs index 91a0bb8..f553f5f 100644 --- a/src/admin/dead_letters.rs +++ b/src/admin/dead_letters.rs @@ -77,6 +77,7 @@ pub struct ListQuery { #[derive(Deserialize)] pub struct CountQuery { + pub collection: Option, pub resolved: Option, } @@ -208,16 +209,26 @@ pub(super) async fn count( _ => "", }; + let collection_clause = if query.collection.is_some() { + " AND collection = ?" + } else { + "" + }; + let mut total: i64 = 0; for table in [ DeadLetterSource::LegacyHooks.table(), DeadLetterSource::Scripts.table(), ] { let sql = adapt_sql( - &format!("SELECT COUNT(*) FROM {table} WHERE 1=1{resolved_clause}"), + &format!("SELECT COUNT(*) FROM {table} WHERE 1=1{resolved_clause}{collection_clause}"), backend, ); - let (n,): (i64,) = sqlx::query_as(&sql) + let mut q = sqlx::query_as::<_, (i64,)>(&sql); + if let Some(ref c) = query.collection { + q = q.bind(c); + } + let (n,) = q .fetch_one(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to count dead letters: {e}")))?; @@ -285,23 +296,15 @@ pub(super) async fn bulk_dismiss( let now = now_rfc3339(); if body.all == Some(true) { - // Operate against both tables. Collection filter only applies - // to legacy rows (the new `dead_letter_scripts` doesn't have a - // collection column — to filter by collection there we'd need - // to JSON-parse `payload`, which is portable-SQL pain. Good - // enough: legacy table is the one with bulk-by-collection - // history anyway.). for source in [DeadLetterSource::LegacyHooks, DeadLetterSource::Scripts] { let table = source.table(); let mut sql = format!("UPDATE {table} SET resolved_at = ? WHERE resolved_at IS NULL"); - let collection_filter = - source == DeadLetterSource::LegacyHooks && body.collection.is_some(); - if collection_filter { + if body.collection.is_some() { sql.push_str(" AND collection = ?"); } let sql = adapt_sql(&sql, backend); let mut q = sqlx::query(&sql).bind(&now); - if collection_filter && let Some(ref c) = body.collection { + if let Some(ref c) = body.collection { q = q.bind(c); } q.execute(&state.db) @@ -444,6 +447,9 @@ async fn list_scripts( "true" => sql.push_str(" AND resolved_at IS NOT NULL"), _ => {} } + if collection.is_some() { + sql.push_str(" AND collection = ?"); + } if cursor.is_some() { sql.push_str(" AND created_at < ?"); } @@ -465,6 +471,9 @@ async fn list_scripts( Option, ), >(&sql); + if let Some(c) = collection { + q = q.bind(c); + } if let Some(cur) = cursor { q = q.bind(cur); } @@ -475,7 +484,7 @@ async fn list_scripts( .await .map_err(|e| AppError::Internal(format!("failed to query scripts dead letters: {e}")))?; - let summaries: Vec = rows + Ok(rows .into_iter() .map(|row| { summary_from_scripts_row(&ScriptsDeadLetterRow { @@ -490,19 +499,7 @@ async fn list_scripts( resolved_at: row.8, }) }) - .collect(); - - // Filter by collection in Rust since the column lives inside - // `payload`. Negligible cost — already client-side after the - // unfiltered DB fetch. - Ok(if let Some(want) = collection { - summaries - .into_iter() - .filter(|s| s.collection == want) - .collect() - } else { - summaries - }) + .collect()) } struct ScriptsDeadLetterRow { @@ -845,18 +842,18 @@ async fn resolve_bulk_ids(state: &AppState, body: &BulkRequest) -> Result(&sql) + { + let mut sql = + String::from("SELECT id FROM dead_letter_scripts WHERE resolved_at IS NULL"); + if body.collection.is_some() { + sql.push_str(" AND collection = ?"); + } + let sql = adapt_sql(&sql, backend); + let mut q = sqlx::query_as::<_, (i64,)>(&sql); + if let Some(ref c) = body.collection { + q = q.bind(c); + } + let rows = q .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to resolve bulk ids: {e}")))?; diff --git a/src/lua/scripts.rs b/src/lua/scripts.rs index f3f8ef6..e05c751 100644 --- a/src/lua/scripts.rs +++ b/src/lua/scripts.rs @@ -366,6 +366,7 @@ pub async fn run_record_event_script( &event_payload, &last_error, MAX_ATTEMPTS, + Some(payload.nsid), ) .await; log_event( @@ -549,6 +550,11 @@ pub async fn run_label_applied_script( } } } + let collection = resolved + .id + .split_once(':') + .map(|(_, suf)| suf) + .filter(|s| *s != "_actor"); write_dead_letter( state, &resolved, @@ -557,6 +563,7 @@ pub async fn run_label_applied_script( &payload, &last_error, MAX_ATTEMPTS, + collection, ) .await; LabelHookOutcome::Continue(original) @@ -763,12 +770,13 @@ async fn write_dead_letter( payload: &Value, error: &str, attempts: u32, + collection: Option<&str>, ) { let payload_str = serde_json::to_string(payload).unwrap_or_else(|_| "{}".to_string()); let sql = adapt_sql( "INSERT INTO dead_letter_scripts - (script_ref, host_kind, host_id, payload, error, attempts, created_at) - VALUES (?, ?, ?, ?, ?, ?, ?)", + (script_ref, host_kind, host_id, payload, error, attempts, created_at, collection) + VALUES (?, ?, ?, ?, ?, ?, ?, ?)", state.db_backend, ); if let Err(e) = sqlx::query(&sql) @@ -779,6 +787,7 @@ async fn write_dead_letter( .bind(error) .bind(attempts as i64) .bind(now_rfc3339()) + .bind(collection) .execute(&state.db) .await { -- 2.51.2