mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-08 08:02:27 +00:00
improve error messages
This commit is contained in:
+1108
-1102
File diff suppressed because it is too large
Load Diff
@@ -45,7 +45,7 @@ pub enum Error {
|
||||
HexErr(#[from] hex::FromHexError),
|
||||
#[error("Migrating database: {0}")]
|
||||
DatabaseMigration(#[from] MigrateError),
|
||||
#[error("{0}")]
|
||||
#[error(transparent)]
|
||||
Anyhow(#[from] anyhow::Error),
|
||||
}
|
||||
|
||||
|
||||
+16
-3
@@ -318,7 +318,12 @@ pub async fn get_path_for_hash<'c>(
|
||||
w_id
|
||||
)
|
||||
.fetch_one(db)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::InternalErr(format!(
|
||||
"querying getting path for hash {hash} in {w_id}: {e}"
|
||||
))
|
||||
})?;
|
||||
Ok(path)
|
||||
}
|
||||
|
||||
@@ -1036,7 +1041,10 @@ pub async fn push<'c>(
|
||||
let premium_workspace =
|
||||
sqlx::query_scalar!("SELECT premium FROM workspace WHERE id = $1", workspace_id)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::InternalErr(format!("fetching if {workspace_id} is premium: {e}"))
|
||||
})?;
|
||||
|
||||
if !premium_workspace && std::env::var("CLOUD_HOSTED").is_ok() {
|
||||
let rate_limiting_queue = sqlx::query_scalar!(
|
||||
@@ -1090,7 +1098,12 @@ pub async fn push<'c>(
|
||||
workspace_id
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::InternalErr(format!(
|
||||
"fetching language for hash {hash} in {workspace_id}: {e}"
|
||||
))
|
||||
})?;
|
||||
(
|
||||
Some(hash.0),
|
||||
Some(path),
|
||||
|
||||
@@ -106,7 +106,7 @@ pub async fn build_oauth_clients(base_url: &str) -> anyhow::Result<AllClients> {
|
||||
match serde_json::from_str::<HashMap<String, OAuthClient>>(&content) {
|
||||
Ok(clients) => clients,
|
||||
Err(e) => {
|
||||
tracing::error!("Error while deserializing oauth.json: {e}");
|
||||
tracing::error!("deserializing oauth.json: {e}");
|
||||
HashMap::new()
|
||||
}
|
||||
}
|
||||
@@ -279,7 +279,8 @@ async fn create_account(
|
||||
payload.refresh_token
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("creating account in {w_id}: {e}")))?;
|
||||
tx.commit().await?;
|
||||
Ok(id.to_string())
|
||||
}
|
||||
|
||||
@@ -11,7 +11,7 @@ use std::str::FromStr;
|
||||
use crate::{
|
||||
audit::{audit_log, ActionKind},
|
||||
db::{UserDB, DB},
|
||||
error::{self, JsonResult, Result},
|
||||
error::{self, Error, JsonResult, Result},
|
||||
jobs::{self, push, JobPayload},
|
||||
users::Authed,
|
||||
utils::{get_owner_from_path, Pagination, StripPath},
|
||||
@@ -149,7 +149,8 @@ async fn create_schedule(
|
||||
ns.enabled
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("inserting schedule in {w_id}: {e}")))?;
|
||||
|
||||
audit_log(
|
||||
&mut tx,
|
||||
@@ -249,7 +250,8 @@ async fn edit_schedule(
|
||||
w_id,
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("updating schedule in {w_id}: {e}")))?;
|
||||
|
||||
if schedule.enabled {
|
||||
tx = push_scheduled_job(tx, schedule).await?;
|
||||
|
||||
@@ -619,7 +619,8 @@ async fn archive_script_by_path(
|
||||
&w_id
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("archiving script in {w_id}: {e}")))?;
|
||||
audit_log(
|
||||
&mut tx,
|
||||
&authed.username,
|
||||
@@ -647,7 +648,9 @@ async fn archive_script_by_hash(
|
||||
)
|
||||
.bind(&hash.0)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("archiving script in {w_id}: {e}")))?;
|
||||
|
||||
audit_log(
|
||||
&mut tx,
|
||||
&authed.username,
|
||||
@@ -679,7 +682,9 @@ async fn delete_script_by_hash(
|
||||
.bind(&hash.0)
|
||||
.bind(&w_id)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("deleting script by hash {w_id}: {e}")))?;
|
||||
|
||||
audit_log(
|
||||
&mut tx,
|
||||
&authed.username,
|
||||
|
||||
@@ -653,7 +653,9 @@ async fn global_whoami(
|
||||
email
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("fetching global identity: {e}")))?;
|
||||
|
||||
Ok(Json(user))
|
||||
}
|
||||
|
||||
@@ -1037,7 +1039,8 @@ async fn set_password(
|
||||
&email
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("setting password: {e}")))?
|
||||
.unwrap_or("".to_string());
|
||||
|
||||
if custom_type != "password".to_string() {
|
||||
|
||||
@@ -43,7 +43,8 @@ pub async fn require_super_admin<'c>(
|
||||
email.as_ref()
|
||||
)
|
||||
.fetch_one(db)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("fetching super admin: {e}")))?;
|
||||
if !is_admin {
|
||||
Err(Error::NotAuthorized(
|
||||
"This endpoint require caller to be a super admin".to_owned(),
|
||||
|
||||
@@ -442,7 +442,9 @@ pub async fn build_crypt<'c>(
|
||||
w_id
|
||||
)
|
||||
.fetch_one(db)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("fetching crypt key: {e}")))?;
|
||||
|
||||
Ok(magic_crypt::new_magic_crypt!(key, 256))
|
||||
}
|
||||
|
||||
|
||||
@@ -927,9 +927,10 @@ async fn handle_child(
|
||||
let canceled = sqlx::query_scalar!("SELECT canceled FROM queue WHERE id = $1", id)
|
||||
.fetch_one(db)
|
||||
.await
|
||||
.map_err(|_| tracing::error!("error getting canceled for id {}", id));
|
||||
.map_err(|e| tracing::error!("error getting canceled for id {}: {e}", id))
|
||||
.unwrap_or(false);
|
||||
|
||||
if canceled.unwrap_or(false) {
|
||||
if canceled {
|
||||
tracing::info!("killed after cancel: {}", job.id);
|
||||
done.store(true, Ordering::Relaxed);
|
||||
}
|
||||
@@ -1008,8 +1009,8 @@ pub async fn restart_zombie_jobs_periodically(
|
||||
mut rx: tokio::sync::broadcast::Receiver<()>,
|
||||
) {
|
||||
loop {
|
||||
let restarted = sqlx::query_scalar!(
|
||||
"UPDATE queue SET running = false WHERE last_ping < $1 RETURNING id",
|
||||
let restarted = sqlx::query!(
|
||||
"UPDATE queue SET running = false WHERE last_ping < $1 and running = true RETURNING id, workspace_id",
|
||||
chrono::Utc::now() - chrono::Duration::seconds(timeout as i64 * 2)
|
||||
)
|
||||
.fetch_all(db)
|
||||
@@ -1017,8 +1018,8 @@ pub async fn restart_zombie_jobs_periodically(
|
||||
.ok()
|
||||
.unwrap_or_else(|| vec![]);
|
||||
|
||||
if restarted.len() > 0 {
|
||||
tracing::info!("restarted zombie jobs {restarted:?}");
|
||||
for r in restarted {
|
||||
tracing::info!("restarted zombie jobs {} {}", r.id, r.workspace_id);
|
||||
}
|
||||
|
||||
tokio::select! {
|
||||
|
||||
@@ -66,7 +66,8 @@ pub async fn update_flow_status_after_job_completion(
|
||||
w_id
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("fetching flow status: {e}")))?
|
||||
.ok_or_else(|| Error::InternalErr(format!("requiring a previous status")))?;
|
||||
|
||||
let old_status = serde_json::from_value::<FlowStatus>(old_status_json)
|
||||
@@ -125,11 +126,7 @@ pub async fn update_flow_status_after_job_completion(
|
||||
.bind(flow)
|
||||
.fetch_one(&mut tx)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::InternalErr(format!(
|
||||
"error during retrieval of stop_early_expr from state: {e}"
|
||||
))
|
||||
})?;
|
||||
.map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e}")))?;
|
||||
|
||||
tracing::debug!("UPDATE: {:?}", new_status);
|
||||
|
||||
@@ -297,7 +294,8 @@ pub async fn get_step_of_flow_status(db: &DB, id: Uuid) -> error::Result<i32> {
|
||||
id
|
||||
)
|
||||
.fetch_one(db)
|
||||
.await?
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("fetching step flow status: {e}")))?
|
||||
.ok_or_else(|| Error::InternalErr(format!("not found step")))?;
|
||||
Ok(r)
|
||||
}
|
||||
|
||||
@@ -191,7 +191,9 @@ async fn get_settings(
|
||||
&w_id
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("getting settings: {e}")))?;
|
||||
|
||||
tx.commit().await?;
|
||||
Ok(Json(settings))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user