mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 08:02:19 +00:00
fix: run custom schedule handlers as the schedule (#11398)
* fix: run custom schedule handlers as the schedule Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: correct handler identity comment and pin preset path check Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: pin preset handler paths instead of matching their name Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: warn on unpinned preset-named schedule handlers Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * chore: update ee-repo-ref to 12a83524b81368043bc939caeb3764fd27a32523 This commit updates the EE repository reference after PR #833 was merged in windmill-ee-private. Previous ee-repo-ref: bcfb0838f105d5bffbe49c6ed787116662ebe796 New ee-repo-ref: 12a83524b81368043bc939caeb3764fd27a32523 Automated by sync-ee-ref workflow. --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com> Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
windmill-internal-app[bot]
parent
6064aecdec
commit
bda00c6cad
@@ -1 +1 @@
|
||||
dfe7b8b9204d21e0264fbea1c6f6eedf9e738d56
|
||||
12a83524b81368043bc939caeb3764fd27a32523
|
||||
|
||||
+16
-1
@@ -4871,6 +4871,18 @@ async fn test_script_schedule_handlers(db: Pool<Postgres>) -> anyhow::Result<()>
|
||||
"a script was run after main job execution but was not schedule error handler"
|
||||
);
|
||||
}
|
||||
|
||||
let (handler_permissioned_as, schedule_permissioned_as): (String, String) =
|
||||
sqlx::query_as(
|
||||
"SELECT j.permissioned_as, s.permissioned_as FROM v2_job j, schedule s
|
||||
WHERE j.id = $1 AND s.workspace_id = j.workspace_id
|
||||
AND s.path = 'f/system/failing_script_schedule'",
|
||||
)
|
||||
.bind(uuid)
|
||||
.fetch_one(&db2)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(handler_permissioned_as, schedule_permissioned_as);
|
||||
},
|
||||
port,
|
||||
)
|
||||
@@ -5786,7 +5798,10 @@ async fn test_scoped_custom_tag_pattern_admission(db: Pool<Postgres>) -> anyhow:
|
||||
"scoped pattern refused elsewhere",
|
||||
!allowed(&db, "other", "cpu-secret", "").await,
|
||||
),
|
||||
("confined tag refused past the patterns it fits", confined.is_err()),
|
||||
(
|
||||
"confined tag refused past the patterns it fits",
|
||||
confined.is_err(),
|
||||
),
|
||||
(
|
||||
"global pattern for the tags nothing confines",
|
||||
allowed(&db, "test-workspace", "gpu-large", "").await,
|
||||
|
||||
@@ -1669,6 +1669,14 @@ async fn run_slack_message_test_job(
|
||||
Path(w_id): Path<String>,
|
||||
Json(req): Json<RunSlackMessageTestJobRequest>,
|
||||
) -> JsonResult<RunSlackMessageTestJobResponse> {
|
||||
// Runs as the error handler identity, handed the Slack bot token.
|
||||
if !(windmill_queue::is_preset_handler_path(&req.hub_script_path)
|
||||
&& req.hub_script_path.ends_with("-slack"))
|
||||
{
|
||||
return Err(Error::BadRequest(
|
||||
"Only a preset Slack handler from the hub can be tested".to_string(),
|
||||
));
|
||||
}
|
||||
let mut fake_result = HashMap::new();
|
||||
fake_result.insert("error".to_string(), to_raw_value(&req.test_msg));
|
||||
fake_result.insert("success_result".to_string(), to_raw_value(&req.test_msg));
|
||||
@@ -1693,7 +1701,7 @@ async fn run_slack_message_test_job(
|
||||
Some(Utc::now()),
|
||||
Some(sqlx::types::Json(to_raw_value(&extra_args))),
|
||||
authed.email.as_str(),
|
||||
false,
|
||||
None,
|
||||
false,
|
||||
None, // Note: we could mark it as high priority to return result quickly to the user
|
||||
)
|
||||
|
||||
@@ -2772,15 +2772,16 @@ async fn send_email_with_instance_smtp(
|
||||
return Err(anyhow::anyhow!("Feature not supported in cloud hosted windmill").into());
|
||||
}
|
||||
|
||||
// Any code pushed as a workspace or schedule error handler, custom ones included, runs as
|
||||
// one of these identities: this keeps out ad-hoc job tokens, not who authors handler code.
|
||||
// Workspace error handlers (admin-configured, custom ones included) and the preset hub
|
||||
// handlers of a schedule run as one of these identities. A custom schedule handler runs as
|
||||
// the schedule and is set by any schedule writer, so it must stay out of this list.
|
||||
let is_handler_job = authed.email == EMAIL_ERROR_HANDLER_USER_EMAIL
|
||||
|| authed.email == ERROR_HANDLER_USER_EMAIL
|
||||
|| authed.email == SCHEDULE_ERROR_HANDLER_USER_EMAIL;
|
||||
|
||||
if !is_handler_job && !windmill_api_auth::is_super_admin_authed(&db, &authed).await? {
|
||||
return Err(Error::NotAuthorized(
|
||||
"Only super admin or a workspace/schedule error handler job can send emails with the instance SMTP"
|
||||
"Only super admin, a workspace error handler or a preset schedule handler can send emails with the instance SMTP"
|
||||
.to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
@@ -163,7 +163,7 @@ const MAX_FREE_EXECS: i32 = 1000;
|
||||
const MAX_FREE_CONCURRENT_RUNS: i32 = 30;
|
||||
|
||||
const ERROR_HANDLER_USERNAME: &str = "error_handler";
|
||||
const SCHEDULE_ERROR_HANDLER_USERNAME: &str = "schedule_error_handler";
|
||||
pub const SCHEDULE_ERROR_HANDLER_USERNAME: &str = "schedule_error_handler";
|
||||
const GLOBAL_ERROR_HANDLER_USERNAME: &str = "global";
|
||||
const SUCCESS_HANDLER_USERNAME: &str = "success_handler";
|
||||
|
||||
@@ -2339,7 +2339,7 @@ pub async fn send_error_to_global_handler<'a, T: Serialize + Send + Sync>(
|
||||
queued_job.started_at,
|
||||
None,
|
||||
&queued_job.permissioned_as_email,
|
||||
false,
|
||||
None,
|
||||
true,
|
||||
None,
|
||||
)
|
||||
@@ -2393,7 +2393,7 @@ pub async fn report_error_to_workspace_handler_or_critical_side_channel(
|
||||
queued_job.started_at,
|
||||
error_handler_extra_args,
|
||||
&queued_job.permissioned_as_email,
|
||||
false,
|
||||
None,
|
||||
false,
|
||||
None,
|
||||
)
|
||||
@@ -2568,7 +2568,7 @@ pub async fn send_error_to_workspace_handler<'a, 'c, T: Serialize + Send + Sync>
|
||||
queued_job.started_at,
|
||||
error_handler_extra_args,
|
||||
&queued_job.permissioned_as_email,
|
||||
false,
|
||||
None,
|
||||
false,
|
||||
None,
|
||||
)
|
||||
@@ -2930,69 +2930,109 @@ pub async fn try_schedule_next_job<'c>(
|
||||
(tx, push_err)
|
||||
}
|
||||
|
||||
pub const ERROR_HANDLER_PATH_TEAMS: &str = "/workspace-or-schedule-error-handler-teams";
|
||||
pub const ERROR_HANDLER_PATH_SLACK: &str = "/workspace-or-schedule-error-handler-slack";
|
||||
pub const ERROR_HANDLER_PATH_EMAIL: &str = "/workspace-or-error-handler-email";
|
||||
const ERROR_HANDLER_PATH_EMAIL: &str = "/workspace-or-error-handler-email";
|
||||
|
||||
enum ErrorHandlerType {
|
||||
Custom,
|
||||
Teams,
|
||||
Slack,
|
||||
Email,
|
||||
/// Every handler version the frontend has offered as a preset (`frontend/src/lib/hubPaths.json`
|
||||
/// and its history, which stored schedules still point at). Matched whole: a hub script is
|
||||
/// resolved by its numeric id alone, so the name after it proves nothing about its code.
|
||||
const PRESET_HANDLER_PATHS: [&str; 23] = [
|
||||
"hub/6512/workspace-or-schedule-error-handler-slack",
|
||||
"hub/5792/workspace-or-schedule-error-handler-slack",
|
||||
"hub/9079/workspace-or-schedule-error-handler-slack",
|
||||
"hub/9206/workspace-or-schedule-error-handler-slack",
|
||||
"hub/19741/workspace-or-schedule-error-handler-slack",
|
||||
"hub/28241/workspace-or-schedule-error-handler-slack",
|
||||
"hub/28794/workspace-or-schedule-error-handler-slack",
|
||||
"hub/2431/slack/schedule-error-handler-slack",
|
||||
"hub/2430/slack/schedule-recovery-handler-slack",
|
||||
"hub/9067/slack/schedule-recovery-handler-slack",
|
||||
"hub/9080/slack/schedule-recovery-handler-slack",
|
||||
"hub/28239/slack/schedule-recovery-handler-slack",
|
||||
"hub/28791/slack/schedule-recovery-handler-slack",
|
||||
"hub/9069/slack/schedule-success-handler-slack",
|
||||
"hub/9072/slack/schedule-success-handler-slack",
|
||||
"hub/28220/slack/schedule-success-handler-slack",
|
||||
"hub/28240/slack/schedule-success-handler-slack",
|
||||
"hub/28793/slack/schedule-success-handler-slack",
|
||||
"hub/11598/workspace-or-schedule-error-handler-teams",
|
||||
"hub/19742/workspace-or-schedule-error-handler-teams",
|
||||
"hub/11593/schedule-recovery-handler-teams",
|
||||
"hub/11596/schedule-success-handler-teams",
|
||||
"hub/19795/workspace-or-error-handler-email",
|
||||
];
|
||||
|
||||
/// The preset handler Windmill ships on the hub that `handler_path` names, `None` for user code.
|
||||
fn preset_handler(handler_path: &str) -> Option<&'static str> {
|
||||
let path = handler_path.strip_prefix("script/").unwrap_or(handler_path);
|
||||
PRESET_HANDLER_PATHS.into_iter().find(|p| *p == path)
|
||||
}
|
||||
|
||||
impl ErrorHandlerType {
|
||||
fn from_error_handler_path(error_handler_path: &str) -> Option<ErrorHandlerType> {
|
||||
let error_handler_path = if error_handler_path.starts_with("script/") {
|
||||
error_handler_path.strip_prefix("script/").unwrap()
|
||||
} else if error_handler_path.starts_with("flow/") {
|
||||
error_handler_path.strip_prefix("flow/").unwrap()
|
||||
} else {
|
||||
error_handler_path
|
||||
};
|
||||
/// Whether `handler_path` is a preset hub handler, which runs as the shared handler identity
|
||||
/// (Slack bot token, instance SMTP) instead of as the schedule or user that set it.
|
||||
pub fn is_preset_handler_path(handler_path: &str) -> bool {
|
||||
preset_handler(handler_path).is_some()
|
||||
}
|
||||
|
||||
if let Some(from_hub) = error_handler_path.strip_prefix("hub/") {
|
||||
let handler_type = if from_hub.ends_with(ERROR_HANDLER_PATH_TEAMS) {
|
||||
ErrorHandlerType::Teams
|
||||
} else if from_hub.ends_with(ERROR_HANDLER_PATH_SLACK) {
|
||||
ErrorHandlerType::Slack
|
||||
} else if from_hub.ends_with(ERROR_HANDLER_PATH_EMAIL) {
|
||||
ErrorHandlerType::Email
|
||||
} else {
|
||||
return None;
|
||||
};
|
||||
|
||||
return Some(handler_type);
|
||||
}
|
||||
|
||||
Some(ErrorHandlerType::Custom)
|
||||
/// `(email, permissioned_as)` a schedule handler runs as. The handler identity reads the
|
||||
/// workspace Slack bot token and may send through the instance SMTP, and a schedule's handlers
|
||||
/// are set by anyone who can write the schedule: only the preset hub handlers get it, custom
|
||||
/// code runs as the schedule itself.
|
||||
pub async fn schedule_handler_identity(
|
||||
db: &Pool<Postgres>,
|
||||
w_id: &str,
|
||||
handler_path: &str,
|
||||
schedule_permissioned_as: &str,
|
||||
preset_email: &str,
|
||||
) -> error::Result<(String, String)> {
|
||||
if is_preset_handler_path(handler_path) {
|
||||
return Ok((
|
||||
preset_email.to_string(),
|
||||
ERROR_HANDLER_USER_GROUP.to_string(),
|
||||
));
|
||||
}
|
||||
if handler_path.contains("hub/") && handler_path.contains("-handler-") {
|
||||
tracing::warn!(
|
||||
"schedule handler {handler_path} in {w_id} is named like a preset but is not a pinned \
|
||||
preset version: it runs as the schedule ({schedule_permissioned_as})"
|
||||
);
|
||||
}
|
||||
let email =
|
||||
windmill_common::users::get_email_from_permissioned_as(schedule_permissioned_as, w_id, db)
|
||||
.await?;
|
||||
Ok((email, schedule_permissioned_as.to_string()))
|
||||
}
|
||||
|
||||
fn get_email_and_permissioned_as(
|
||||
async fn error_handler_identity(
|
||||
db: &Pool<Postgres>,
|
||||
w_id: &str,
|
||||
error_handler_path: &str,
|
||||
is_global_error_handler: bool,
|
||||
is_schedule_error_handler: bool,
|
||||
) -> (&'static str, String) {
|
||||
let res = if is_global_error_handler {
|
||||
(SUPERADMIN_SECRET_EMAIL, SUPERADMIN_SECRET_EMAIL.to_string())
|
||||
} else if is_schedule_error_handler {
|
||||
(
|
||||
schedule_permissioned_as: Option<&str>,
|
||||
) -> error::Result<(String, String)> {
|
||||
if is_global_error_handler {
|
||||
return Ok((
|
||||
SUPERADMIN_SECRET_EMAIL.to_string(),
|
||||
SUPERADMIN_SECRET_EMAIL.to_string(),
|
||||
));
|
||||
}
|
||||
if let Some(schedule_permissioned_as) = schedule_permissioned_as {
|
||||
return schedule_handler_identity(
|
||||
db,
|
||||
w_id,
|
||||
error_handler_path,
|
||||
schedule_permissioned_as,
|
||||
SCHEDULE_ERROR_HANDLER_USER_EMAIL,
|
||||
ERROR_HANDLER_USER_GROUP.to_string(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
let email = if preset_handler(error_handler_path)
|
||||
.is_some_and(|p| p.ends_with(ERROR_HANDLER_PATH_EMAIL))
|
||||
{
|
||||
EMAIL_ERROR_HANDLER_USER_EMAIL
|
||||
} else {
|
||||
let handler_type = ErrorHandlerType::from_error_handler_path(error_handler_path);
|
||||
|
||||
let email = match handler_type {
|
||||
Some(ErrorHandlerType::Email) => EMAIL_ERROR_HANDLER_USER_EMAIL,
|
||||
_ => ERROR_HANDLER_USER_EMAIL,
|
||||
};
|
||||
|
||||
(email, ERROR_HANDLER_USER_GROUP.to_string())
|
||||
ERROR_HANDLER_USER_EMAIL
|
||||
};
|
||||
|
||||
res
|
||||
Ok((email.to_string(), ERROR_HANDLER_USER_GROUP.to_string()))
|
||||
}
|
||||
|
||||
pub async fn push_error_handler<'a, 'c, T: Serialize + Send + Sync>(
|
||||
@@ -3008,7 +3048,8 @@ pub async fn push_error_handler<'a, 'c, T: Serialize + Send + Sync>(
|
||||
started_at: Option<DateTime<Utc>>,
|
||||
extra_args: Option<Json<Box<RawValue>>>,
|
||||
email: &str,
|
||||
is_schedule_error_handler: bool,
|
||||
// The schedule's `permissioned_as` when this is a schedule's own error handler.
|
||||
schedule_permissioned_as: Option<&str>,
|
||||
is_global_error_handler: bool,
|
||||
priority: Option<i16>,
|
||||
) -> windmill_common::error::Result<Uuid> {
|
||||
@@ -3047,17 +3088,17 @@ pub async fn push_error_handler<'a, 'c, T: Serialize + Send + Sync>(
|
||||
|
||||
let result = sanitize_result(result);
|
||||
|
||||
let (email, permissioned_as) = if let Some(on_behalf_of) = on_behalf_of.as_ref() {
|
||||
(
|
||||
on_behalf_of.email.as_str(),
|
||||
on_behalf_of.permissioned_as.clone(),
|
||||
)
|
||||
let (email, permissioned_as) = if let Some(on_behalf_of) = on_behalf_of {
|
||||
(on_behalf_of.email, on_behalf_of.permissioned_as)
|
||||
} else {
|
||||
get_email_and_permissioned_as(
|
||||
error_handler_identity(
|
||||
db,
|
||||
w_id,
|
||||
on_failure_path,
|
||||
is_global_error_handler,
|
||||
is_schedule_error_handler,
|
||||
schedule_permissioned_as,
|
||||
)
|
||||
.await?
|
||||
};
|
||||
|
||||
let tx = PushIsolationLevel::IsolatedRoot(db.clone());
|
||||
@@ -3069,12 +3110,12 @@ pub async fn push_error_handler<'a, 'c, T: Serialize + Send + Sync>(
|
||||
PushArgs { extra: Some(extra), args: &result },
|
||||
if is_global_error_handler {
|
||||
GLOBAL_ERROR_HANDLER_USERNAME
|
||||
} else if is_schedule_error_handler {
|
||||
} else if schedule_permissioned_as.is_some() {
|
||||
SCHEDULE_ERROR_HANDLER_USERNAME
|
||||
} else {
|
||||
ERROR_HANDLER_USERNAME
|
||||
},
|
||||
email,
|
||||
&email,
|
||||
permissioned_as,
|
||||
Some(&format!("error.handler.{job_id}")),
|
||||
None,
|
||||
@@ -8145,6 +8186,41 @@ mod git_sync_concurrency_key_tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod preset_handler_path_tests {
|
||||
use super::is_preset_handler_path;
|
||||
|
||||
#[test]
|
||||
fn only_pinned_hub_scripts_are_presets() {
|
||||
assert!(is_preset_handler_path(
|
||||
"script/hub/28791/slack/schedule-recovery-handler-slack"
|
||||
));
|
||||
// The hub resolves by id: any other id is arbitrary hub code whatever its name.
|
||||
assert!(!is_preset_handler_path(
|
||||
"hub/13968/workspace-or-schedule-error-handler-slack"
|
||||
));
|
||||
assert!(!is_preset_handler_path(
|
||||
"script/u/me/workspace-or-schedule-error-handler-slack"
|
||||
));
|
||||
}
|
||||
|
||||
/// A preset version bumped in the frontend without being pinned here would silently run as
|
||||
/// the schedule's owner and lose access to the Slack bot token.
|
||||
#[test]
|
||||
fn frontend_presets_are_pinned() {
|
||||
let hub_paths: std::collections::HashMap<String, String> =
|
||||
serde_json::from_str(include_str!("../../../frontend/src/lib/hubPaths.json")).unwrap();
|
||||
for (name, path) in hub_paths {
|
||||
if path.contains("-handler-") {
|
||||
assert!(
|
||||
is_preset_handler_path(&path),
|
||||
"{name}: {path} is not pinned"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod result_metadata_tests {
|
||||
use super::{ResultMetadata, ValidableJson};
|
||||
|
||||
Reference in New Issue
Block a user