fix: retry on inserting completed job (#4784)

* fix: retry on pushing next scheduled job of schedule

* all

* improve error handling

* update sqlx
This commit is contained in:
Ruben Fiszel
2024-11-25 01:33:44 +01:00
committed by GitHub
parent 5ed0ae697b
commit 6768e5bbbd
12 changed files with 332 additions and 280 deletions
@@ -5,7 +5,7 @@
"columns": [
{
"ordinal": 0,
"name": "bool",
"name": "?column?",
"type_info": "Bool"
}
],
@@ -5,7 +5,7 @@
"columns": [
{
"ordinal": 0,
"name": "bool",
"name": "?column?",
"type_info": "Bool"
}
],
@@ -5,7 +5,7 @@
"columns": [
{
"ordinal": 0,
"name": "bool",
"name": "?column?",
"type_info": "Bool"
}
],
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO completed_job AS cj\n ( workspace_id\n , id\n , parent_job\n , created_by\n , created_at\n , started_at\n , duration_ms\n , success\n , script_hash\n , script_path\n , args\n , result\n , raw_code\n , raw_lock\n , canceled\n , canceled_by\n , canceled_reason\n , job_kind\n , schedule_path\n , permissioned_as\n , flow_status\n , raw_flow\n , is_flow_step\n , is_skipped\n , language\n , email\n , visible_to_owner\n , mem_peak\n , tag\n , priority\n )\n VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000, $7, $8, $9,$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29)\n ON CONFLICT (id) DO UPDATE SET success = $7, result = $11 RETURNING duration_ms",
"query": "INSERT INTO completed_job AS cj\n ( workspace_id\n , id\n , parent_job\n , created_by\n , created_at\n , started_at\n , duration_ms\n , success\n , script_hash\n , script_path\n , args\n , result\n , raw_code\n , raw_lock\n , canceled\n , canceled_by\n , canceled_reason\n , job_kind\n , schedule_path\n , permissioned_as\n , flow_status\n , raw_flow\n , is_flow_step\n , is_skipped\n , language\n , email\n , visible_to_owner\n , mem_peak\n , tag\n , priority\n )\n VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), COALESCE($30::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000), $7, $8, $9,$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29)\n ON CONFLICT (id) DO UPDATE SET success = $7, result = $11 RETURNING duration_ms",
"describe": {
"columns": [
{
@@ -87,12 +87,13 @@
"Bool",
"Int4",
"Varchar",
"Int2"
"Int2",
"Int8"
]
},
"nullable": [
false
]
},
"hash": "d5a8614286c170e0d175903cd1b53ff66b37ed8110a0b67aedb9f25e6a7383e1"
"hash": "f9fc0084fe086ef80005bb64a8bb6b493e53583017c18e2ab44f44125c52d548"
}
+1
View File
@@ -10916,6 +10916,7 @@ version = "1.430.1"
dependencies = [
"anyhow",
"async-recursion",
"backon",
"base64 0.22.1",
"bit-vec",
"bytes",
+270 -246
View File
@@ -39,7 +39,6 @@ use windmill_audit::audit_ee::{audit_log, AuditAuthor};
use windmill_audit::ActionKind;
use windmill_common::{
add_time,
auth::{fetch_authed_from_permissioned_as, permissioned_as_to_username},
db::{Authed, UserDB},
error::{self, to_anyhow, Error},
@@ -65,6 +64,9 @@ use windmill_common::{
DB, METRICS_ENABLED,
};
use backon::ConstantBuilder;
use backon::{BackoffBuilder, Retryable};
#[cfg(feature = "enterprise")]
use windmill_common::BASE_URL;
@@ -173,8 +175,7 @@ pub async fn cancel_single_job<'c>(
e,
"server",
false,
#[cfg(feature = "benchmark")]
&mut windmill_common::bench::BenchmarkIter::new(),
None,
)
.await;
@@ -462,7 +463,7 @@ pub async fn add_completed_job_error(
e: serde_json::Value,
_worker_name: &str,
flow_is_done: bool,
#[cfg(feature = "benchmark")] bench: &mut windmill_common::bench::BenchmarkIter,
duration: Option<i64>,
) -> Result<WrappedError, Error> {
#[cfg(feature = "prometheus")]
register_metric(
@@ -499,8 +500,7 @@ pub async fn add_completed_job_error(
mem_peak,
canceled_by,
flow_is_done,
#[cfg(feature = "benchmark")]
bench,
duration,
)
.await?;
Ok(result)
@@ -520,134 +520,139 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
mem_peak: i32,
canceled_by: Option<CanceledBy>,
flow_is_done: bool,
#[cfg(feature = "benchmark")] bench: &mut windmill_common::bench::BenchmarkIter,
duration: Option<i64>,
) -> Result<Uuid, Error> {
// tracing::error!("Start");
// let start = tokio::time::Instant::now();
add_time!(bench, "add_completed_job start");
// add_time!(bench, "add_completed_job start");
if !result.is_valid_json() {
return Err(Error::InternalErr(
"Result of job is invalid json (empty)".to_string(),
));
}
let mut tx = db.begin().await?;
let _job_id = queued_job.id;
let (opt_uuid, _duration, _skip_downstream_error_handlers) = (|| async {
let mut tx = db.begin().await?;
let job_id = queued_job.id;
// tracing::error!("1 {:?}", start.elapsed());
let job_id = queued_job.id;
// tracing::error!("1 {:?}", start.elapsed());
tracing::debug!(
"completed job {} {}",
queued_job.id,
serde_json::to_string(&result).unwrap_or_else(|_| "".to_string())
);
tracing::debug!(
"completed job {} {}",
queued_job.id,
serde_json::to_string(&result).unwrap_or_else(|_| "".to_string())
);
let (raw_code, raw_lock, raw_flow) = if !*MIN_VERSION_IS_AT_LEAST_1_427.read().await {
sqlx::query!(
"SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json<Box<JsonRawValue>>\"
FROM job WHERE id = $1 AND workspace_id = $2 LIMIT 1",
&job_id,
&queued_job.workspace_id
)
.fetch_one(db)
.map_ok(|record| (record.raw_code, record.raw_lock, record.raw_flow))
.or_else(|_| {
let (raw_code, raw_lock, raw_flow) = if !*MIN_VERSION_IS_AT_LEAST_1_427.read().await {
sqlx::query!(
"SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json<Box<JsonRawValue>>\"
FROM queue WHERE id = $1 AND workspace_id = $2 LIMIT 1",
FROM job WHERE id = $1 AND workspace_id = $2 LIMIT 1",
&job_id,
&queued_job.workspace_id
)
.fetch_one(db)
.map_ok(|record| (record.raw_code, record.raw_lock, record.raw_flow))
})
.await
.unwrap_or_default()
} else {
(None, None, None)
};
let mem_peak = mem_peak.max(queued_job.mem_peak.unwrap_or(0));
add_time!(bench, "add_completed_job query START");
let _duration: i64 = sqlx::query_scalar!(
"INSERT INTO completed_job AS cj
( workspace_id
, id
, parent_job
, created_by
, created_at
, started_at
, duration_ms
, success
, script_hash
, script_path
, args
, result
, raw_code
, raw_lock
, canceled
, canceled_by
, canceled_reason
, job_kind
, schedule_path
, permissioned_as
, flow_status
, raw_flow
, is_flow_step
, is_skipped
, language
, email
, visible_to_owner
, mem_peak
, tag
, priority
.or_else(|_| {
sqlx::query!(
"SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json<Box<JsonRawValue>>\"
FROM queue WHERE id = $1 AND workspace_id = $2 LIMIT 1",
&job_id,
&queued_job.workspace_id
)
VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000, $7, $8, $9,\
$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29)
ON CONFLICT (id) DO UPDATE SET success = $7, result = $11 RETURNING duration_ms",
queued_job.workspace_id,
queued_job.id,
queued_job.parent_job,
queued_job.created_by,
queued_job.created_at,
queued_job.started_at,
success,
queued_job.script_hash.map(|x| x.0),
queued_job.script_path,
&queued_job.args as &Option<Json<HashMap<String, Box<RawValue>>>>,
result as Json<&T>,
raw_code,
raw_lock,
canceled_by.is_some(),
canceled_by.clone().map(|cb| cb.username).flatten(),
canceled_by.clone().map(|cb| cb.reason).flatten(),
queued_job.job_kind.clone() as JobKind,
queued_job.schedule_path,
queued_job.permissioned_as,
&queued_job.flow_status as &Option<Json<Box<RawValue>>>,
&raw_flow as &Option<Json<Box<RawValue>>>,
queued_job.is_flow_step,
skipped,
queued_job.language.clone() as Option<ScriptLang>,
queued_job.email,
queued_job.visible_to_owner,
if mem_peak > 0 { Some(mem_peak) } else { None },
queued_job.tag,
queued_job.priority,
)
.fetch_one(&mut *tx)
.await
.map_err(|e| Error::InternalErr(format!("Could not add completed job {job_id}: {e:#}")))?;
// tracing::error!("2 {:?}", start.elapsed());
.fetch_one(db)
.map_ok(|record| (record.raw_code, record.raw_lock, record.raw_flow))
})
.await
.unwrap_or_default()
} else {
(None, None, None)
};
add_time!(bench, "add_completed_job query END");
let mem_peak = mem_peak.max(queued_job.mem_peak.unwrap_or(0));
// add_time!(bench, "add_completed_job query START");
if !queued_job.is_flow_step {
if _duration > 500
&& (queued_job.job_kind == JobKind::Script || queued_job.job_kind == JobKind::Preview)
{
if let Err(e) = sqlx::query!(
let _duration = sqlx::query_scalar!(
"INSERT INTO completed_job AS cj
( workspace_id
, id
, parent_job
, created_by
, created_at
, started_at
, duration_ms
, success
, script_hash
, script_path
, args
, result
, raw_code
, raw_lock
, canceled
, canceled_by
, canceled_reason
, job_kind
, schedule_path
, permissioned_as
, flow_status
, raw_flow
, is_flow_step
, is_skipped
, language
, email
, visible_to_owner
, mem_peak
, tag
, priority
)
VALUES ($1, $2, $3, $4, $5, COALESCE($6, now()), COALESCE($30::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($6, now()))))*1000), $7, $8, $9,\
$10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29)
ON CONFLICT (id) DO UPDATE SET success = $7, result = $11 RETURNING duration_ms",
queued_job.workspace_id,
queued_job.id,
queued_job.parent_job,
queued_job.created_by,
queued_job.created_at,
queued_job.started_at,
success,
queued_job.script_hash.map(|x| x.0),
queued_job.script_path,
&queued_job.args as &Option<Json<HashMap<String, Box<RawValue>>>>,
result as Json<&T>,
raw_code,
raw_lock,
canceled_by.is_some(),
canceled_by.clone().map(|cb| cb.username).flatten(),
canceled_by.clone().map(|cb| cb.reason).flatten(),
queued_job.job_kind.clone() as JobKind,
queued_job.schedule_path,
queued_job.permissioned_as,
&queued_job.flow_status as &Option<Json<Box<RawValue>>>,
&raw_flow as &Option<Json<Box<RawValue>>>,
queued_job.is_flow_step,
skipped,
queued_job.language.clone() as Option<ScriptLang>,
queued_job.email,
queued_job.visible_to_owner,
if mem_peak > 0 { Some(mem_peak) } else { None },
queued_job.tag,
queued_job.priority,
duration,
)
.fetch_one(&mut *tx)
.await
.map_err(|e| Error::InternalErr(format!("Could not add completed job {job_id}: {e:#}")))?;
// tracing::error!("2 {:?}", start.elapsed());
// add_time!(bench, "add_completed_job query END");
if !queued_job.is_flow_step {
if _duration > 500
&& (queued_job.job_kind == JobKind::Script
|| queued_job.job_kind == JobKind::Preview)
{
if let Err(e) = sqlx::query!(
"UPDATE completed_job SET flow_status = q.flow_status FROM queue q WHERE completed_job.id = $1 AND q.id = $1 AND q.workspace_id = $2 AND completed_job.workspace_id = $2 AND q.flow_status IS NOT NULL",
&queued_job.id,
&queued_job.workspace_id
@@ -656,9 +661,9 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
.await {
tracing::error!("Could not update job duration: {}", e);
}
}
if let Some(parent_job) = queued_job.parent_job {
if let Err(e) = sqlx::query_scalar!(
}
if let Some(parent_job) = queued_job.parent_job {
if let Err(e) = sqlx::query_scalar!(
"UPDATE queue SET flow_status = jsonb_set(jsonb_set(COALESCE(flow_status, '{}'::jsonb), array[$1], COALESCE(flow_status->$1, '{}'::jsonb)), array[$1, 'duration_ms'], to_jsonb($2::bigint)) WHERE id = $3 AND workspace_id = $4",
&queued_job.id.to_string(),
_duration,
@@ -669,106 +674,106 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
.await {
tracing::error!("Could not update parent job flow_status: {}", e);
}
}
}
}
// tracing::error!("Added completed job {:#?}", queued_job);
#[cfg(feature = "enterprise")]
let mut skip_downstream_error_handlers = false;
tx = delete_job(tx, &queued_job.workspace_id, job_id).await?;
// tracing::error!("3 {:?}", start.elapsed());
// tracing::error!("Added completed job {:#?}", queued_job);
if queued_job.is_flow_step {
if let Some(parent_job) = queued_job.parent_job {
// persist the flow last progress timestamp to avoid zombie flow jobs
tracing::debug!(
"Persisting flow last progress timestamp to flow job: {:?}",
parent_job
);
sqlx::query!(
let mut _skip_downstream_error_handlers = false;
tx = delete_job(tx, &queued_job.workspace_id, job_id).await?;
// tracing::error!("3 {:?}", start.elapsed());
if queued_job.is_flow_step {
if let Some(parent_job) = queued_job.parent_job {
// persist the flow last progress timestamp to avoid zombie flow jobs
tracing::debug!(
"Persisting flow last progress timestamp to flow job: {:?}",
parent_job
);
sqlx::query!(
"UPDATE queue SET last_ping = now() WHERE id = $1 AND workspace_id = $2 AND canceled = false",
parent_job,
&queued_job.workspace_id
)
.execute(&mut *tx)
.await?;
if flow_is_done {
let r = sqlx::query_scalar!(
if flow_is_done {
let r = sqlx::query_scalar!(
"UPDATE parallel_monitor_lock SET last_ping = now() WHERE parent_flow_id = $1 and job_id = $2 RETURNING 1",
parent_job,
&queued_job.id
).fetch_optional(&mut *tx).await?;
if r.is_some() {
tracing::info!(
if r.is_some() {
tracing::info!(
"parallel flow iteration is done, setting parallel monitor last ping lock for job {}",
&queued_job.id
);
}
}
}
}
} else {
if queued_job.schedule_path.is_some() && queued_job.script_path.is_some() {
let schedule_path = queued_job.schedule_path.as_ref().unwrap();
let script_path = queued_job.script_path.as_ref().unwrap();
} else {
if queued_job.schedule_path.is_some() && queued_job.script_path.is_some() {
let schedule_path = queued_job.schedule_path.as_ref().unwrap();
let script_path = queued_job.script_path.as_ref().unwrap();
let schedule =
get_schedule_opt(&mut tx, &queued_job.workspace_id, schedule_path).await?;
let schedule =
get_schedule_opt(&mut tx, &queued_job.workspace_id, schedule_path).await?;
if let Some(schedule) = schedule {
#[cfg(feature = "enterprise")]
{
skip_downstream_error_handlers = schedule.ws_error_handler_muted;
}
if let Some(schedule) = schedule {
#[cfg(feature = "enterprise")]
{
_skip_downstream_error_handlers = schedule.ws_error_handler_muted;
}
// script or flow that failed on start and might not have been rescheduled
let schedule_next_tick = !queued_job.is_flow()
|| {
let flow_status = queued_job.parse_flow_status();
flow_status.is_some_and(|fs| {
// script or flow that failed on start and might not have been rescheduled
let schedule_next_tick = !queued_job.is_flow()
|| {
let flow_status = queued_job.parse_flow_status();
flow_status.is_some_and(|fs| {
fs.step == 0
&& fs.modules.get(0).is_some_and(|m| {
matches!(m, FlowStatusModule::WaitingForPriorSteps { .. }) || matches!(m, FlowStatusModule::Failure { job, ..} if job == &Uuid::nil())
})
})
};
};
if schedule_next_tick {
if let Err(err) = handle_maybe_scheduled_job(
if schedule_next_tick {
if let Err(err) = handle_maybe_scheduled_job(
db,
queued_job,
&schedule,
script_path,
&queued_job.workspace_id,
)
.await
{
match err {
Error::QuotaExceeded(_) => (),
// scheduling next job failed and could not disable schedule => make zombie job to retry
_ => return Ok((Some(job_id), 0, true)),
}
};
}
#[cfg(feature = "enterprise")]
if let Err(err) = apply_schedule_handlers(
db,
queued_job,
&schedule,
script_path,
&queued_job.workspace_id,
success,
result,
job_id,
queued_job.started_at.unwrap_or(chrono::Utc::now()),
queued_job.priority,
)
.await
{
match err {
Error::QuotaExceeded(_) => return Err(err.into()),
// scheduling next job failed and could not disable schedule => make zombie job to retry
_ => return Ok(job_id),
}
};
}
#[cfg(feature = "enterprise")]
if let Err(err) = apply_schedule_handlers(
db,
&schedule,
script_path,
&queued_job.workspace_id,
success,
result,
job_id,
queued_job.started_at.unwrap_or(chrono::Utc::now()),
queued_job.priority,
)
.await
{
if !success {
tracing::error!("Could not apply schedule error handler: {}", err);
let base_url = BASE_URL.read().await;
let w_id: &String = &queued_job.workspace_id;
if !matches!(err, Error::QuotaExceeded(_)) {
report_error_to_workspace_handler_or_critical_side_channel(
if !success {
tracing::error!("Could not apply schedule error handler: {}", err);
let base_url = BASE_URL.read().await;
let w_id: &String = &queued_job.workspace_id;
if !matches!(err, Error::QuotaExceeded(_)) {
report_error_to_workspace_handler_or_critical_side_channel(
&queued_job,
db,
format!(
@@ -778,82 +783,104 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
),
)
.await;
}
} else {
tracing::error!("Could not apply schedule recovery handler: {}", err);
}
} else {
tracing::error!("Could not apply schedule recovery handler: {}", err);
}
};
} else {
tracing::error!(
};
} else {
tracing::error!(
"Schedule {schedule_path} in {} not found. Impossible to schedule again and apply schedule handlers",
&queued_job.workspace_id
);
}
}
}
}
if queued_job.concurrent_limit.is_some() {
let concurrency_key = match concurrency_key(db, queued_job).await {
Ok(c) => c,
Err(e) => {
tracing::error!(
"Could not get concurrency key for job {} defaulting to default key: {e:?}",
queued_job.id
);
legacy_concurrency_key(db, queued_job)
.await
.unwrap_or_else(|| queued_job.full_path_with_workspace())
}
};
if let Err(e) = sqlx::query_scalar!(
if queued_job.concurrent_limit.is_some() {
let concurrency_key = match concurrency_key(db, queued_job).await {
Ok(c) => c,
Err(e) => {
tracing::error!(
"Could not get concurrency key for job {} defaulting to default key: {e:?}",
queued_job.id
);
legacy_concurrency_key(db, queued_job)
.await
.unwrap_or_else(|| queued_job.full_path_with_workspace())
}
};
if let Err(e) = sqlx::query_scalar!(
"UPDATE concurrency_counter SET job_uuids = job_uuids - $2 WHERE concurrency_id = $1",
concurrency_key,
queued_job.id.hyphenated().to_string(),
)
.execute(&mut *tx)
.await
{
tracing::error!("Could not decrement concurrency counter: {}", e);
}
.execute(&mut *tx)
.await
{
tracing::error!("Could not decrement concurrency counter: {}", e);
}
if let Err(e) = sqlx::query_scalar!(
"UPDATE concurrency_key SET ended_at = now() WHERE job_id = $1",
queued_job.id,
)
.execute(&mut *tx)
.await
.map_err(|e| {
Error::InternalErr(format!(
if let Err(e) = sqlx::query_scalar!(
"UPDATE concurrency_key SET ended_at = now() WHERE job_id = $1",
queued_job.id,
)
.execute(&mut *tx)
.await
.map_err(|e| {
Error::InternalErr(format!(
"Error updating to add ended_at timestamp concurrency_key={concurrency_key}: {e:#}"
))
}) {
tracing::error!("Could not update concurrency_key: {}", e);
}) {
tracing::error!("Could not update concurrency_key: {}", e);
}
tracing::debug!("decremented concurrency counter");
}
tracing::debug!("decremented concurrency counter");
}
if JOB_TOKEN.is_none() {
sqlx::query!("DELETE FROM job_perms WHERE job_id = $1", job_id)
.execute(&mut *tx)
.await?;
}
if JOB_TOKEN.is_none() {
sqlx::query!("DELETE FROM job_perms WHERE job_id = $1", job_id)
.execute(&mut *tx)
.await?;
}
tx.commit().await?;
tracing::info!(
%job_id,
root_job = ?queued_job.root_job.map(|x| x.to_string()).unwrap_or_else(|| String::new()),
path = &queued_job.script_path(),
job_kind = ?queued_job.job_kind,
started_at = ?queued_job.started_at.map(|x| x.to_string()).unwrap_or_else(|| String::new()),
duration = ?_duration,
permissioned_as = ?queued_job.permissioned_as,
email = ?queued_job.email,
created_by = queued_job.created_by,
is_flow_step = queued_job.is_flow_step,
language = ?queued_job.language,
success,
"inserted completed job: {} (success: {success})",
queued_job.id
);
tx.commit().await?;
tracing::info!(
%job_id,
root_job = ?queued_job.root_job.map(|x| x.to_string()).unwrap_or_else(|| String::new()),
path = &queued_job.script_path(),
job_kind = ?queued_job.job_kind,
started_at = ?queued_job.started_at.map(|x| x.to_string()).unwrap_or_else(|| String::new()),
duration = ?_duration,
permissioned_as = ?queued_job.permissioned_as,
email = ?queued_job.email,
created_by = queued_job.created_by,
is_flow_step = queued_job.is_flow_step,
language = ?queued_job.language,
success,
"inserted completed job: {} (success: {success})",
queued_job.id
);
Ok((None, _duration, _skip_downstream_error_handlers)) as windmill_common::error::Result<(Option<Uuid>, i64, bool)>
})
.retry(
ConstantBuilder::default()
.with_delay(std::time::Duration::from_secs(3))
.with_max_times(5)
.build(),
)
.when(|err| !matches!(err, Error::QuotaExceeded(_)))
.notify(|err, dur| {
tracing::error!(
"Could not insert completed job, retrying in {dur:#?}, err: {err:#?}"
);
})
.sleep(tokio::time::sleep)
.await?;
// if scheduling next job failed, return the job_id early to ensure the job get retried after a timeout
if let Some(job_id) = opt_uuid {
return Ok(job_id);
}
#[cfg(feature = "cloud")]
if *CLOUD_HOSTED && !queued_job.is_flow() && _duration > 1000 {
@@ -938,10 +965,10 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
),
)
.await;
} else if !skip_downstream_error_handlers
} else if !_skip_downstream_error_handlers
&& (matches!(queued_job.job_kind, JobKind::Script)
|| matches!(queued_job.job_kind, JobKind::Flow)
&& !has_failure_module(db, job_id).await)
&& !has_failure_module(db, _job_id).await)
&& queued_job.parent_job.is_none()
{
let result = serde_json::from_str(
@@ -1232,9 +1259,6 @@ pub async fn send_error_to_workspace_handler<'a, 'c, T: Serialize + Send + Sync>
Ok(())
}
use backon::ConstantBuilder;
use backon::{BackoffBuilder, Retryable};
#[instrument(level = "trace", skip_all)]
pub async fn handle_maybe_scheduled_job<'c>(
db: &Pool<Postgres>,
+1 -1
View File
@@ -90,7 +90,7 @@ object_store = { workspace = true, optional = true}
convert_case.workspace = true
yaml-rust.workspace = true
swc_ecma_parser.workspace = true
backon.workspace = true
[build-dependencies]
deno_fetch = { workspace = true, optional = true }
@@ -185,14 +185,14 @@ pub async fn handle_dedicated_process(
let result = Arc::new(result);
append_logs(&job.id, &job.workspace_id, logs.clone(), db).await;
if line.starts_with("wm_res[success]:") {
job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: true, cached_res_path: None, token: token.to_string() }).await.unwrap()
job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: true, cached_res_path: None, token: token.to_string(), duration: None }).await.unwrap()
} else {
job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string() }).await.unwrap()
job_completed_tx.send(JobCompleted { job , result, mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string(), duration: None }).await.unwrap()
}
},
Err(e) => {
tracing::error!("Could not deserialize job result `{line}`: {e:?}");
job_completed_tx.send(JobCompleted { job , result: Arc::new(to_raw_value(&serde_json::json!({"error": format!("Could not deserialize job result `{line}`: {e:?}")}))), mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string() }).await.unwrap();
job_completed_tx.send(JobCompleted { job , result: Arc::new(to_raw_value(&serde_json::json!({"error": format!("Could not deserialize job result `{line}`: {e:?}")}))), mem_peak: 0, canceled_by: None, success: false, cached_res_path: None, token: token.to_string(), duration: None }).await.unwrap();
},
};
logs = init_log.clone();
+19 -10
View File
@@ -186,8 +186,18 @@ async fn send_job_completed(
success: bool,
cached_res_path: Option<String>,
token: String,
duration: Option<i64>,
) {
let jc = JobCompleted { job, result, mem_peak, canceled_by, success, cached_res_path, token };
let jc = JobCompleted {
job,
result,
mem_peak,
canceled_by,
success,
cached_res_path,
token,
duration,
};
job_completed_tx.send(jc).await.expect("send job completed")
}
@@ -203,6 +213,7 @@ pub async fn process_result(
column_order: Option<Vec<String>>,
new_args: Option<HashMap<String, Box<RawValue>>>,
db: &DB,
duration: Option<i64>,
) -> error::Result<bool> {
match result {
Ok(r) => {
@@ -257,6 +268,7 @@ pub async fn process_result(
true,
cached_res_path,
token,
duration,
)
.await;
Ok(true)
@@ -300,6 +312,7 @@ pub async fn process_result(
false,
cached_res_path,
token,
duration,
)
.await;
Ok(false)
@@ -362,7 +375,7 @@ pub async fn handle_receive_completed_job(
#[tracing::instrument(name = "completed_job", level = "info", skip_all, fields(job_id = %job.id))]
pub async fn process_completed_job(
JobCompleted { job, result, mem_peak, success, cached_res_path, canceled_by, .. }: JobCompleted,
JobCompleted { job, result, mem_peak, success, cached_res_path, canceled_by, duration, .. }: JobCompleted,
client: &AuthedClient,
db: &DB,
worker_dir: &str,
@@ -391,8 +404,7 @@ pub async fn process_completed_job(
mem_peak.to_owned(),
canceled_by,
false,
#[cfg(feature = "benchmark")]
bench,
duration,
)
.await?;
drop(job);
@@ -434,8 +446,7 @@ pub async fn process_completed_job(
),
worker_name,
false,
#[cfg(feature = "benchmark")]
bench,
None,
)
.await?;
if job.is_flow_step {
@@ -501,8 +512,7 @@ pub async fn handle_job_error(
err.clone(),
worker_name,
false,
#[cfg(feature = "benchmark")]
bench,
None,
)
.await
};
@@ -562,8 +572,7 @@ pub async fn handle_job_error(
e,
worker_name,
false,
#[cfg(feature = "benchmark")]
bench,
None,
)
.await;
}
+22 -2
View File
@@ -107,6 +107,9 @@ use crate::{
},
};
use backon::ConstantBuilder;
use backon::{BackoffBuilder, Retryable};
#[cfg(feature = "enterprise")]
use crate::dedicated_worker::create_dedicated_worker_map;
@@ -1107,7 +1110,7 @@ pub async fn run_worker(
let (occupancy_rate, occupancy_rate_15s, occupancy_rate_5m, occupancy_rate_30m) =
occupancy_metrics.update_occupancy_metrics();
if let Err(e) = sqlx::query!(
if let Err(e) = (|| sqlx::query!(
"UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2,
occupancy_rate = $3, memory_usage = $4, wm_memory_usage = $5, vcpus = COALESCE($7, vcpus),
memory = COALESCE($8, memory), occupancy_rate_15s = $9, occupancy_rate_5m = $10, occupancy_rate_30m = $11 WHERE worker = $6",
@@ -1122,7 +1125,19 @@ pub async fn run_worker(
occupancy_rate_15s,
occupancy_rate_5m,
occupancy_rate_30m
).execute(db).await {
).execute(db)).retry(
ConstantBuilder::default()
.with_delay(std::time::Duration::from_secs(2))
.with_max_times(10)
.build(),
)
.notify(|err, dur| {
tracing::error!(
"retrying updating worker ping in {dur:#?}, err: {err:#?}"
);
})
.sleep(tokio::time::sleep)
.await {
tracing::error!("failed to update worker ping, exiting: {}", e);
killpill_tx.send(()).unwrap_or_default();
}
@@ -1349,6 +1364,7 @@ pub async fn run_worker(
cached_res_path: None,
token: "".to_string(),
canceled_by: None,
duration: None,
})
.await
.expect("send job completed END");
@@ -1704,6 +1720,7 @@ pub struct JobCompleted {
pub cached_res_path: Option<String>,
pub token: String,
pub canceled_by: Option<CanceledBy>,
pub duration: Option<i64>,
}
async fn do_nativets(
@@ -1828,6 +1845,7 @@ async fn handle_queued_job(
None
};
let started = Instant::now();
let (raw_code, raw_lock, raw_flow) = match (raw_code, raw_lock, raw_flow) {
(None, None, None) => sqlx::query!(
"SELECT raw_code, raw_lock, raw_flow AS \"raw_flow: Json<Box<RawValue>>\"
@@ -1922,6 +1940,7 @@ async fn handle_queued_job(
success: true,
cached_res_path: None,
token: authed_client.token,
duration: None,
})
.await
.expect("send job completed");
@@ -2087,6 +2106,7 @@ async fn handle_queued_job(
column_order,
new_args,
db,
Some(started.elapsed().as_millis() as i64),
)
.await
}
+6 -10
View File
@@ -1033,8 +1033,7 @@ pub async fn update_flow_status_after_job_completion_internal(
canceled_job_to_result(&flow_job),
worker_name,
true,
#[cfg(feature = "benchmark")]
bench,
None,
)
.await?;
} else {
@@ -1074,9 +1073,8 @@ pub async fn update_flow_status_after_job_completion_internal(
Json(&nresult),
0,
None,
true,
#[cfg(feature = "benchmark")]
bench,
true,
None,
)
.await?;
} else {
@@ -1092,9 +1090,8 @@ pub async fn update_flow_status_after_job_completion_internal(
),
0,
None,
true,
#[cfg(feature = "benchmark")]
bench,
true,
None,
)
.await?;
}
@@ -1131,8 +1128,7 @@ pub async fn update_flow_status_after_job_completion_internal(
e,
worker_name,
true,
#[cfg(feature = "benchmark")]
bench,
None,
)
.await;
true
+3 -2
View File
@@ -210,10 +210,11 @@
<div class="text-red-600 text-xs"> Invalid cron syntax </div>
{/if}
</svelte:fragment>
<div class="flex flex-row-reverse text-2xs text-tertiary -mt-1">
<div class="flex flex-row-reverse text-2xs text-tertiary -mt-1 hover:underline">
<a
class="text-tertiary"
href="https://www.windmill.dev/docs/core_concepts/scheduling#cron-syntax">Croner</a
href="https://www.windmill.dev/docs/core_concepts/scheduling#cron-syntax"
target="_blank">Croner</a
>
</div>
<input