add more duration warnings to logs

This commit is contained in:
Ruben Fiszel
2025-05-14 10:13:14 +02:00
parent d5fbeaf91b
commit ef068ab5ef
2 changed files with 67 additions and 13 deletions
+6 -2
View File
@@ -719,7 +719,10 @@ impl<F: Future> Future for WarnAfterFuture<F> {
// Poll the timeout future to check if it has elapsed.
if !*this.warned {
if this.timeout.poll(cx).is_ready() {
tracing::warn!(location = this.location, "SLOW_QUERY: query to db taking longer than expected (> {} seconds). This is a sign the database is under heavy load, query is too heavy or database is undersized",
tracing::warn!(
location = this.location,
"SLOW_QUERY: query {} to db taking longer than expected (> {} seconds)",
this.location,
this.seconds,
);
*this.warned = true;
@@ -733,7 +736,8 @@ impl<F: Future> Future for WarnAfterFuture<F> {
let elapsed = this.start_time.elapsed();
tracing::warn!(
location = this.location,
"SLOW_QUERY: completed with total duration: {:.2?}",
"SLOW_QUERY: completed query {} with total duration: {:.2?}",
this.location,
elapsed
);
}
+61 -11
View File
@@ -1665,6 +1665,7 @@ async fn push_next_flow_job(
worker_dir: worker_dir.to_string(),
token: client.token.clone(),
})
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::internal_err(format!(
@@ -1684,6 +1685,7 @@ async fn push_next_flow_job(
flow_job.workspace_id.as_str()
)
.fetch_one(db)
.warn_after_seconds(3)
.await?;
if no_flow_overlap {
let overlapping = sqlx::query_scalar!(
@@ -1704,6 +1706,7 @@ async fn push_next_flow_job(
flow_job.runnable_path()
)
.fetch_all(db)
.warn_after_seconds(3)
.await?;
if overlapping.len() > 0 {
let overlapping_str = overlapping
@@ -1724,6 +1727,7 @@ async fn push_next_flow_job(
worker_dir: worker_dir.to_string(),
token: client.token.clone(),
})
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::internal_err(format!(
@@ -1749,6 +1753,7 @@ async fn push_next_flow_job(
flow_job.scheduled_for.to_string(),
)]),
)
.warn_after_seconds(3)
.await?;
if skip {
job_completed_tx
@@ -1761,6 +1766,7 @@ async fn push_next_flow_job(
worker_dir: worker_dir.to_string(),
token: client.token.clone(),
})
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::internal_err(format!(
@@ -1787,7 +1793,10 @@ async fn push_next_flow_job(
if last_job_result.is_some() {
last_job_result.unwrap()
} else {
match get_previous_job_result(db, flow_job.workspace_id.as_str(), &status).await? {
match get_previous_job_result(db, flow_job.workspace_id.as_str(), &status)
.warn_after_seconds(3)
.await?
{
None => Arc::new(to_raw_value(&json!("{}"))),
Some(previous_job_result) => Arc::new(previous_job_result),
}
@@ -1806,7 +1815,7 @@ async fn push_next_flow_job(
FlowStatusModule::WaitingForPriorSteps { .. } | FlowStatusModule::WaitingForEvents { .. }
) {
if let Some((suspend, last)) = needs_resume(&flow, &status) {
let mut tx = db.begin().await?;
let mut tx = db.begin().warn_after_seconds(3).await?;
/* Lock this row to prevent the suspend column getting out out of sync
* if a resume message arrives after we fetch and count them here.
@@ -1817,6 +1826,7 @@ async fn push_next_flow_job(
flow_job.id
)
.fetch_one(&mut *tx)
.warn_after_seconds(3)
.await
.context("lock flow in queue")?;
@@ -1825,7 +1835,9 @@ async fn push_next_flow_job(
)
.bind(last)
.fetch_all(&mut *tx)
.await?
.warn_after_seconds(3)
.await
?
.into_iter()
.collect::<Vec<_>>();
@@ -1865,6 +1877,7 @@ async fn push_next_flow_job(
None,
None
)
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::ExecutionErr(format!(
@@ -1897,6 +1910,7 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
}
@@ -1938,6 +1952,7 @@ async fn push_next_flow_job(
Some(&serde_json::json!({"approved": false, "job_id": flow_job.id, "details": "Suspend timed out without approval but can continue".to_string()}).to_string()),
None,
)
.warn_after_seconds(3)
.await?;
}
@@ -1959,6 +1974,7 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
// Remove the approval conditions from the flow status
@@ -1969,10 +1985,11 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
/* continue on and run this job! */
tx.commit().await?;
tx.commit().warn_after_seconds(3).await?;
/* not enough messages to do this job, "park"/suspend until there are */
} else if matches!(
@@ -2002,6 +2019,7 @@ async fn push_next_flow_job(
flow_job.id,
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
sqlx::query!(
@@ -2010,9 +2028,10 @@ async fn push_next_flow_job(
flow_job.id,
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
tx.commit().await?;
tx.commit().warn_after_seconds(3).await?;
return Ok(None);
/* cancelled or we're WaitingForEvents but we don't have enough messages (timed out) */
@@ -2027,9 +2046,10 @@ async fn push_next_flow_job(
Some(&serde_json::json!({"approved": false, "job_id": flow_job.id, "details": "Suspend timed out without approval and is cancelled".to_string()}).to_string()),
None,
)
.warn_after_seconds(3)
.await?;
}
tx.commit().await?;
tx.commit().warn_after_seconds(3).await?;
let (logs, error_name) = if let Some(disapprover) = is_disapproved {
(
@@ -2057,6 +2077,7 @@ async fn push_next_flow_job(
logs.clone(),
&db.into(),
)
.warn_after_seconds(3)
.await;
job_completed_tx
@@ -2069,6 +2090,7 @@ async fn push_next_flow_job(
worker_dir: worker_dir.to_string(),
token: client.token.clone(),
})
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::internal_err(format!(
@@ -2137,6 +2159,7 @@ async fn push_next_flow_job(
None,
None,
)
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::ExecutionErr(format!(
@@ -2237,6 +2260,7 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(db)
.warn_after_seconds(3)
.await
.context("update flow retry")?;
};
@@ -2255,7 +2279,9 @@ async fn push_next_flow_job(
drop(resume_messages);
let is_skipped = if let Some(skip_if) = &module.skip_if {
let idcontext = get_transform_context(&flow_job, previous_id.as_str(), &status).await?;
let idcontext = get_transform_context(&flow_job, previous_id.as_str(), &status)
.warn_after_seconds(3)
.await?;
compute_bool_from_expr(
&skip_if.expr,
arc_flow_job_args.clone(),
@@ -2266,6 +2292,7 @@ async fn push_next_flow_job(
Some((resumes.clone(), resume.clone(), approvers.clone())),
None,
)
.warn_after_seconds(3)
.await?
} else {
false
@@ -2295,6 +2322,7 @@ async fn push_next_flow_job(
&flow_job.workspace_id
)
.fetch_optional(db)
.warn_after_seconds(3)
.await?;
if let Some(args) = args {
Ok(Marc::new(args.map(|x| x.0).unwrap_or_else(HashMap::new)))
@@ -2327,7 +2355,9 @@ async fn push_next_flow_job(
| FlowModuleValue::FlowScript { input_transforms, .. }
| FlowModuleValue::Flow { input_transforms, .. },
) => {
let ctx = get_transform_context(&flow_job, &previous_id, &status).await?;
let ctx = get_transform_context(&flow_job, &previous_id, &status)
.warn_after_seconds(3)
.await?;
transform_context = Some(ctx);
let by_id = transform_context.as_ref().unwrap();
transform_input(
@@ -2340,6 +2370,7 @@ async fn push_next_flow_job(
by_id,
client,
)
.warn_after_seconds(3)
.await
.map(Marc::new)
}
@@ -2370,6 +2401,7 @@ async fn push_next_flow_job(
approvers.clone(),
is_skipped,
)
.warn_after_seconds(3)
.await?;
tracing::info!(id = %flow_job.id, root_id = %job_root, "next flow transform computed");
@@ -2395,6 +2427,7 @@ async fn push_next_flow_job(
flow_job.id
)
.fetch_optional(db)
.warn_after_seconds(3)
.await?
.flatten();
@@ -2433,7 +2466,7 @@ async fn push_next_flow_job(
};
let len = job_payloads.len();
let mut tx = db.begin().await?;
let mut tx = db.begin().warn_after_seconds(3).await?;
let nargs = args.as_ref();
for (i, payload_tag) in job_payloads.into_iter().enumerate() {
if i % 100 == 0 && i != 0 {
@@ -2443,6 +2476,7 @@ async fn push_next_flow_job(
flow_job.id,
)
.execute(db)
.warn_after_seconds(3)
.await?;
}
tracing::debug!(id = %flow_job.id, root_id = %job_root, "pushing job {i} of {len}");
@@ -2481,7 +2515,9 @@ async fn push_next_flow_job(
args.insert("iter".to_string(), to_raw_value(new_args));
if let Some(input_transforms) = simple_input_transforms {
//previous id is none because we do not want to use previous id if we are in a for loop
let ctx = get_transform_context(&flow_job, "", &status).await?;
let ctx = get_transform_context(&flow_job, "", &status)
.warn_after_seconds(3)
.await?;
let ti = transform_input(
Marc::new(args),
arc_last_job_result.clone(),
@@ -2492,6 +2528,7 @@ async fn push_next_flow_job(
&ctx,
client,
)
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::ExecutionErr(
@@ -2527,7 +2564,9 @@ async fn push_next_flow_job(
to_raw_value(&json!({ "index": i as i32, "value": itered[i]})),
);
if let Some(input_transforms) = simple_input_transforms {
let ctx = get_transform_context(&flow_job, &previous_id, &status).await?;
let ctx = get_transform_context(&flow_job, &previous_id, &status)
.warn_after_seconds(3)
.await?;
let ti = transform_input(
Marc::new(hm),
arc_last_job_result.clone(),
@@ -2538,6 +2577,7 @@ async fn push_next_flow_job(
&ctx,
client,
)
.warn_after_seconds(3)
.await
.map_err(|e| {
Error::ExecutionErr(format!(
@@ -2608,6 +2648,7 @@ async fn push_next_flow_job(
flow_job.workspace_id,
)
.fetch_optional(&mut *tx)
.warn_after_seconds(3)
.await?
.map(|x| x.into())
} else {
@@ -2668,6 +2709,7 @@ async fn push_next_flow_job(
worker_name
)
.execute(&mut *inner_tx)
.warn_after_seconds(3)
.await;
}
@@ -2688,6 +2730,7 @@ async fn push_next_flow_job(
uuid,
)
.execute(&mut *inner_tx)
.warn_after_seconds(3)
.await?;
}
tracing::debug!(id = %flow_job.id, root_id = %job_root, "updated suspend for {uuid}");
@@ -2707,6 +2750,7 @@ async fn push_next_flow_job(
root_job.unwrap_or(flow_job.id)
)
.execute(&mut *inner_tx)
.warn_after_seconds(3)
.await?;
}
@@ -2731,6 +2775,7 @@ async fn push_next_flow_job(
uuid
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
tracing::debug!(id = %flow_job.id, root_id = %job_root, "updated parallel monitor lock for {uuid}");
}
@@ -2844,6 +2889,7 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
}
Step::PreprocessorStep => {
@@ -2860,6 +2906,7 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
}
Step::Step(i) => {
@@ -2877,6 +2924,7 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
}
};
@@ -2888,6 +2936,7 @@ async fn push_next_flow_job(
flow_job.id
)
.execute(&mut *tx)
.warn_after_seconds(3)
.await?;
if continue_on_same_worker {
@@ -2903,6 +2952,7 @@ async fn push_next_flow_job(
if continue_on_same_worker {
same_worker_tx
.send(SameWorkerPayload { job_id: first_uuid, recoverable: true })
.warn_after_seconds(3)
.await
.map_err(to_anyhow)?;
}