fix(backend): email triggers error handler and retry (#6601)

This commit is contained in:
hugocasa
2025-09-12 20:51:25 +02:00
committed by GitHub
parent 0e82962545
commit dfd857e57c
5 changed files with 17 additions and 4 deletions

View File

@@ -1 +1 @@
f860277e8a5719c59d8961856fd84e3be5c3a83b
8d3c5d3bdf03b43d6f4b4ede30ab276354874724

View File

@@ -989,6 +989,7 @@ async fn route_job(
trigger.error_handler_path.as_deref(),
trigger.error_handler_args.as_ref(),
format!("http_trigger/{}", trigger.path),
None,
)
.await
.map_err(|e| e.into_response())
@@ -1005,6 +1006,7 @@ async fn route_job(
trigger.error_handler_path.as_deref(),
trigger.error_handler_args.as_ref(),
format!("http_trigger/{}", trigger.path),
None,
)
.await
.map_err(|e| e.into_response())

View File

@@ -500,6 +500,7 @@ pub trait Listener: TriggerCrud + TriggerJobArgs {
error_handler_path.as_deref(),
error_handler_args,
format!("{}_trigger/{}", Self::TRIGGER_KIND, listening_trigger.path),
None,
)
.await?;

View File

@@ -487,6 +487,7 @@ async fn trigger_runnable_inner(
error_handler_path: Option<&str>,
error_handler_args: Option<&sqlx::types::Json<HashMap<String, serde_json::Value>>>,
trigger_path: String,
job_id: Option<Uuid>,
) -> Result<(Uuid, Option<bool>, Option<String>)> {
let error_handler_args = error_handler_args.map(|args| {
let args = args
@@ -499,7 +500,7 @@ async fn trigger_runnable_inner(
let user_db = user_db.unwrap_or_else(|| UserDB::new(db.clone()));
let (uuid, delete_after_use, early_return) = if is_flow {
let run_query = RunJobQuery::default();
let run_query = RunJobQuery { job_id, ..Default::default() };
let path = StripPath(runnable_path.to_string());
let (uuid, early_return) = run_flow_by_path_inner(
authed,
@@ -524,6 +525,7 @@ async fn trigger_runnable_inner(
error_handler_path,
error_handler_args.as_ref(),
trigger_path,
job_id,
)
.await?;
(uuid, delete_after_use, None)
@@ -545,6 +547,7 @@ pub async fn trigger_runnable(
error_handler_path: Option<&str>,
error_handler_args: Option<&sqlx::types::Json<HashMap<String, serde_json::Value>>>,
trigger_path: String,
job_id: Option<Uuid>,
) -> Result<axum::response::Response> {
let (uuid, _, _) = trigger_runnable_inner(
db,
@@ -558,6 +561,7 @@ pub async fn trigger_runnable(
error_handler_path,
error_handler_args,
trigger_path,
job_id,
)
.await?;
Ok((StatusCode::CREATED, uuid.to_string()).into_response())
@@ -590,6 +594,7 @@ pub async fn trigger_runnable_and_wait_for_result(
error_handler_path,
error_handler_args,
trigger_path,
None,
)
.await?;
let (result, success) =
@@ -630,6 +635,7 @@ pub async fn trigger_runnable_and_wait_for_raw_result(
error_handler_path,
error_handler_args,
trigger_path,
None,
)
.await?;
@@ -670,9 +676,10 @@ async fn trigger_script_internal(
error_handler_path: Option<&str>,
error_handler_args: Option<&sqlx::types::Json<HashMap<String, Box<RawValue>>>>,
trigger_path: String,
job_id: Option<Uuid>,
) -> Result<(Uuid, Option<bool>)> {
if retry.is_none() && error_handler_path.is_none() {
let run_query = RunJobQuery::default();
let run_query = RunJobQuery { job_id, ..Default::default() };
let path = StripPath(script_path.to_string());
run_script_by_path_inner(
authed,
@@ -696,6 +703,7 @@ async fn trigger_script_internal(
error_handler_path,
error_handler_args,
trigger_path,
job_id,
)
.await
}
@@ -712,6 +720,7 @@ async fn trigger_script_with_retry_and_error_handler(
error_handler_path: Option<&str>,
error_handler_args: Option<&sqlx::types::Json<HashMap<String, Box<RawValue>>>>,
trigger_path: String,
job_id: Option<Uuid>,
) -> Result<(Uuid, Option<bool>)> {
#[cfg(feature = "enterprise")]
check_license_key_valid().await?;
@@ -805,7 +814,7 @@ async fn trigger_script_with_retry_and_error_handler(
None,
None,
None,
None,
job_id,
false,
false,
None,

View File

@@ -450,6 +450,7 @@ impl Listener for WebsocketTrigger {
error_handler_path,
error_handler_args,
format!("websocket_trigger/{}", listening_trigger.path),
None,
)
.await?;
}