simplify warn_after_seconds

This commit is contained in:
Ruben Fiszel
2024-12-04 18:24:31 +01:00
parent f6c4951877
commit 715d09c30d
4 changed files with 23 additions and 20 deletions

View File

@@ -34,7 +34,7 @@ pub const GIT_VERSION: &str =
git_version!(args = ["--tag", "--always"], fallback = "unknown-version");
use crate::CRITICAL_ALERT_MUTE_UI_ENABLED;
use std::panic::{self, AssertUnwindSafe};
use std::panic::{self, AssertUnwindSafe, Location};
use std::sync::atomic::Ordering;
use crate::worker::CLOUD_HOSTED;
@@ -544,13 +544,16 @@ use pin_project_lite::pin_project;
pub trait WarnAfterExt: Future + Sized {
/// Warns if the future takes longer than the specified number of seconds to complete.
fn warn_after_seconds(self, seconds: u8, location: &'static str) -> WarnAfterFuture<Self> {
#[track_caller]
fn warn_after_seconds(self, seconds: u8) -> WarnAfterFuture<Self> {
let caller = Location::caller();
let location = format!("{}:{}", caller.file(), caller.line());
WarnAfterFuture {
future: self,
timeout: time::sleep(Duration::from_secs(seconds as u64)),
warned: false,
start_time: std::time::Instant::now(),
location,
location: location,
seconds,
}
}
@@ -567,7 +570,7 @@ pin_project! {
#[pin]
timeout: Sleep,
warned: bool,
location: &'static str,
location: String,
start_time: std::time::Instant,
seconds: u8,
}

View File

@@ -432,7 +432,7 @@ pub async fn process_completed_job(
#[cfg(feature = "benchmark")]
bench,
)
.warn_after_seconds(10, "update_flow_status_after_job_completion success")
.warn_after_seconds(10)
.await?;
}
}
@@ -471,7 +471,7 @@ pub async fn process_completed_job(
#[cfg(feature = "benchmark")]
bench,
)
.warn_after_seconds(10, "update_flow_status_after_job_completion error")
.warn_after_seconds(10)
.await?;
}
}

View File

@@ -160,7 +160,7 @@ pub async fn create_token_for_owner_in_bg(
&email,
&job_id,
)
.warn_after_seconds(5, "creating token for owner")
.warn_after_seconds(5)
.await
.expect("could not create job token");
*locked = token;
@@ -1826,7 +1826,7 @@ async fn handle_queued_job(
let daily_count = sqlx::query!(
"SELECT value FROM metrics WHERE id = 'email_trigger_usage' AND created_at > NOW() - INTERVAL '1 day' ORDER BY created_at DESC LIMIT 1"
).fetch_optional(db)
.warn_after_seconds(5, "getting email_trigger_usage")
.warn_after_seconds(5)
.await?.map(|x| serde_json::from_value::<i64>(x.value).unwrap_or(1));
if let Some(count) = daily_count {
@@ -1840,7 +1840,7 @@ async fn handle_queued_job(
serde_json::json!(count + 1)
)
.execute(db)
.warn_after_seconds(5, "updating email_trigger_usage")
.warn_after_seconds(5)
.await?;
}
} else {
@@ -1848,7 +1848,7 @@ async fn handle_queued_job(
"INSERT INTO metrics (id, value) VALUES ('email_trigger_usage', to_jsonb(1))"
)
.execute(db)
.warn_after_seconds(5, "inserting email_trigger_usage")
.warn_after_seconds(5)
.await?;
}
}
@@ -1861,7 +1861,7 @@ async fn handle_queued_job(
.ok_or_else(|| Error::InternalErr(format!("expected parent job")))?,
job.id,
)
.warn_after_seconds(5, "updating flow status in progress")
.warn_after_seconds(5)
.await?;
Some(r)
@@ -1874,7 +1874,7 @@ async fn handle_queued_job(
&job.workspace_id
)
.execute(db)
.warn_after_seconds(5, "updating parent job started_at flow_status")
.warn_after_seconds(5)
.await {
tracing::error!("Could not update parent job started_at flow_status: {}", e);
}
@@ -1891,7 +1891,7 @@ async fn handle_queued_job(
job.workspace_id
)
.fetch_one(db)
.warn_after_seconds(5, "getting job raw values")
.warn_after_seconds(5)
.await
.map(|record| (record.raw_code, record.raw_lock, record.raw_flow))
.unwrap_or_default(),
@@ -1932,7 +1932,7 @@ async fn handle_queued_job(
&job.parent_job.unwrap()
)
.fetch_one(db)
.warn_after_seconds(5, "getting script path from queue for caching purposes")
.warn_after_seconds(5)
.await
.map_err(|e| {
Error::InternalErr(format!(
@@ -1967,7 +1967,7 @@ async fn handle_queued_job(
&job.workspace_id,
&cached_res_path,
)
.warn_after_seconds(5, "getting cached resource value")
.warn_after_seconds(5)
.await;
if let Some(cached_resource_value) = cached_resource_value_maybe {
{
@@ -2006,7 +2006,7 @@ async fn handle_queued_job(
worker_dir,
job_completed_tx.0.clone(),
)
.warn_after_seconds(10, "handling flow")
.warn_after_seconds(10)
.await?;
Ok(true)
} else {

View File

@@ -1110,7 +1110,7 @@ pub async fn update_flow_status_after_job_completion_internal(
worker_dir,
job_completed_tx,
)
.warn_after_seconds(10, "handle_flow in update_flow_status")
.warn_after_seconds(10)
.await
{
Err(err) => {
@@ -1512,7 +1512,7 @@ pub async fn handle_flow(
let schedule_path = flow_job.schedule_path.as_ref().unwrap();
let schedule =
get_schedule_opt(&mut tx, &flow_job.workspace_id, schedule_path).warn_after_seconds(5, "get schedule_opt in handle_flow").await?;
get_schedule_opt(&mut tx, &flow_job.workspace_id, schedule_path).warn_after_seconds(5).await?;
tx.commit().await?;
@@ -1524,7 +1524,7 @@ pub async fn handle_flow(
flow_job.script_path.as_ref().unwrap(),
&flow_job.workspace_id,
)
.warn_after_seconds(5, "handle_maybe_scheduled_job in handle_flow")
.warn_after_seconds(5)
.await
{
match err {
@@ -1552,7 +1552,7 @@ pub async fn handle_flow(
worker_dir,
job_completed_tx,
)
.warn_after_seconds(10, "push next flow job")
.warn_after_seconds(10)
.await?;
Ok(())
}