fix rare stack overflow bc of async size

This commit is contained in:
Ruben Fiszel
2025-10-16 20:05:09 +00:00
parent c2b22b3df8
commit 51c89573bb
2 changed files with 198 additions and 194 deletions
+197 -193
View File
@@ -3838,7 +3838,8 @@ pub async fn push<'c, 'd>(
if let Some(skip_handler) = skip_handler {
let mut skip_input_transforms = HashMap::<String, InputTransform>::new();
for (arg_name, arg_value) in skip_handler.args {
skip_input_transforms.insert(arg_name, InputTransform::Static { value: arg_value });
skip_input_transforms
.insert(arg_name, InputTransform::Static { value: arg_value });
}
modules.push(FlowModule {
@@ -3973,7 +3974,7 @@ pub async fn push<'c, 'd>(
// this is a new flow being pushed, flow_status is set to flow_value:
let flow_status: FlowStatus = FlowStatus::new(&flow_value);
(
None, // No version needed - flow is stored in raw_flow like FlowPreview
None, // No version needed - flow is stored in raw_flow like FlowPreview
Some(path),
None,
JobKind::SingleStepFlow,
@@ -5138,127 +5139,128 @@ pub async fn preprocess_dependency_job(job: &mut PulledJob, db: &DB) -> error::R
let kind = job.kind;
// Handle dependency job debouncing cleanup when a job is pulled for execution
if kind.is_dependency() && !*WMDEBUG_NO_DJOB_DEBOUNCING {
// Only used for testing in tests/relative_imports.rs
// Give us some space to work with.
#[cfg(debug_assertions)]
if let Some(duration) = job
.args
.as_ref()
.map(|x| {
x.get("dbg_sleep_between_pull_and_debounce_key_removal")
.map(|v| serde_json::from_str::<u32>(v.get()).ok())
.flatten()
})
.flatten()
{
tracing::debug!("going to sleep",);
sleep(std::time::Duration::from_secs(duration as u64)).await;
}
tracing::debug!(
"Processing debounce cleanup for dependency job {} at path {:?}",
&job.id,
&job.runnable_path
);
let key = format!("{}:{}:dependency", &job.workspace_id, job.runnable_path());
let mut tx = db.begin().await?;
// === DEBOUNCE CLEANUP ===
//
// Clean up the debounce_key entry for this job (if it exists).
//
// IMPORTANT: We delete by key (not job_id) to avoid race conditions:
// If pusher has locked this row then this call will be blocked until all txs are commited.
//
// The idea is that the worker_lockfiles::trigger_dependents_to_recompute_locks will fetch the latest version of the obj.
// This object needs to be created before the djob is executed and it happens right here.
//
// This way the next pusher can fetch the latest version of object and base their djob payload on newest version.
// The concurrency limit on djobs will make sure that by the time next djob is started executing the base version it is referencing
// has already calculated all locks. This way even next djob will always use the fully finalized version of object.
//
//
//
// Note: We don't use a transaction here for performance (it's called during job pull).
// This means there's a tiny window where the job is running but key isn't deleted yet,
// which is acceptable because new requests will just accumulate data to this job.
tracing::debug!(
job_id = %job.id,
"Cleaning up debounce_key entry for completed/pulled job"
);
// This will either:
// 1. Block until pusher pushed. Which gives us:
// - If there was any stale data in pusher, then we will read it here (couple of lines below)
// 2. Block pusher until we are done here. This gives us:
// - We will clone objects and retrieve the latest version. So when we are done the pusher can read latest version.
sqlx::query!("DELETE FROM debounce_key WHERE key = $1", &key)
.execute(&mut *tx)
.await
.map_err(|e| {
tracing::error!(
error = %e,
job_id = %job.id,
"Failed to delete debounce_key"
);
e
})?;
if job
.args
.as_ref()
.map(|x| x.get("triggered_by_relative_import").is_some())
.unwrap_or_default()
{
let Some(base_hash) = job.runnable_id else {
return Err(Error::InternalErr(
"Missing runnable_id for dependency job triggered by relative import"
.to_string(),
));
};
return Box::pin(async move {
// Only used for testing in tests/relative_imports.rs
// Give us some space to work with.
#[cfg(debug_assertions)]
if let Some(duration) = job
.args
.as_ref()
.map(|x| {
x.get("dbg_sleep_between_pull_and_debounce_key_removal")
.map(|v| serde_json::from_str::<u32>(v.get()).ok())
.flatten()
})
.flatten()
{
tracing::debug!("going to sleep",);
sleep(std::time::Duration::from_secs(duration as u64)).await;
}
tracing::debug!(
job_id = %job.id,
base_hash = %base_hash,
job_kind = ?kind,
"Creating new version for dependency job triggered by relative import"
"Processing debounce cleanup for dependency job {} at path {:?}",
&job.id,
&job.runnable_path
);
let new_id = match kind {
JobKind::Dependencies => {
let deployment_message = job
.args
.clone()
.map(|hashmap| {
hashmap
.get("deployment_message")
.map(|map_value| {
serde_json::from_str::<String>(map_value.get()).ok()
})
.flatten()
})
.flatten();
let key = format!("{}:{}:dependency", &job.workspace_id, job.runnable_path());
let mut tx = db.begin().await?;
// This way we tell downstream which script we should archive when the resolution is finished.
// (not used at the moment)
job.args
.as_mut()
.map(|args| args.insert("base_hash".to_owned(), to_raw_value(&*base_hash)));
// === DEBOUNCE CLEANUP ===
//
// Clean up the debounce_key entry for this job (if it exists).
//
// IMPORTANT: We delete by key (not job_id) to avoid race conditions:
// If pusher has locked this row then this call will be blocked until all txs are commited.
//
// The idea is that the worker_lockfiles::trigger_dependents_to_recompute_locks will fetch the latest version of the obj.
// This object needs to be created before the djob is executed and it happens right here.
//
// This way the next pusher can fetch the latest version of object and base their djob payload on newest version.
// The concurrency limit on djobs will make sure that by the time next djob is started executing the base version it is referencing
// has already calculated all locks. This way even next djob will always use the fully finalized version of object.
//
//
//
// Note: We don't use a transaction here for performance (it's called during job pull).
// This means there's a tiny window where the job is running but key isn't deleted yet,
// which is acceptable because new requests will just accumulate data to this job.
tracing::debug!(
job_id = %job.id,
"Cleaning up debounce_key entry for completed/pulled job"
);
let new_hash = windmill_common::scripts::clone_script(
base_hash,
&job.workspace_id,
deployment_message,
&mut tx,
)
.await?;
// This will either:
// 1. Block until pusher pushed. Which gives us:
// - If there was any stale data in pusher, then we will read it here (couple of lines below)
// 2. Block pusher until we are done here. This gives us:
// - We will clone objects and retrieve the latest version. So when we are done the pusher can read latest version.
sqlx::query!("DELETE FROM debounce_key WHERE key = $1", &key)
.execute(&mut *tx)
.await
.map_err(|e| {
tracing::error!(
error = %e,
job_id = %job.id,
"Failed to delete debounce_key"
);
e
})?;
new_hash
}
JobKind::FlowDependencies => {
sqlx::query_scalar!(
"INSERT INTO flow_version
if job
.args
.as_ref()
.map(|x| x.get("triggered_by_relative_import").is_some())
.unwrap_or_default()
{
let Some(base_hash) = job.runnable_id else {
return Err(Error::InternalErr(
"Missing runnable_id for dependency job triggered by relative import"
.to_string(),
));
};
tracing::debug!(
job_id = %job.id,
base_hash = %base_hash,
job_kind = ?kind,
"Creating new version for dependency job triggered by relative import"
);
let new_id = match kind {
JobKind::Dependencies => {
let deployment_message = job
.args
.clone()
.map(|hashmap| {
hashmap
.get("deployment_message")
.map(|map_value| {
serde_json::from_str::<String>(map_value.get()).ok()
})
.flatten()
})
.flatten();
// This way we tell downstream which script we should archive when the resolution is finished.
// (not used at the moment)
job.args.as_mut().map(|args| {
args.insert("base_hash".to_owned(), to_raw_value(&*base_hash))
});
let new_hash = windmill_common::scripts::clone_script(
base_hash,
&job.workspace_id,
deployment_message,
&mut tx,
)
.await?;
new_hash
}
JobKind::FlowDependencies => {
sqlx::query_scalar!(
"INSERT INTO flow_version
(workspace_id, path, value, schema, created_by)
SELECT workspace_id, path, value, schema, created_by
@@ -5266,98 +5268,100 @@ pub async fn preprocess_dependency_job(job: &mut PulledJob, db: &DB) -> error::R
RETURNING id
",
job.runnable_path(),
job.workspace_id,
*base_hash,
)
.fetch_one(&mut *tx)
.await?
}
JobKind::AppDependencies => {
sqlx::query_scalar!(
"INSERT INTO app_version
job.runnable_path(),
job.workspace_id,
*base_hash,
)
.fetch_one(&mut *tx)
.await?
}
JobKind::AppDependencies => {
sqlx::query_scalar!(
"INSERT INTO app_version
(app_id, value, created_by, raw_app)
SELECT app_id, value, created_by, raw_app
FROM app_version WHERE id = $1
RETURNING id",
*base_hash
)
.fetch_one(&mut *tx)
.await?
}
_ => {
return Err(Error::InternalErr(format!(
"Matched unexpected JobKind ({:?}). This is a bug!",
kind
)))
}
};
*base_hash
)
.fetch_one(&mut *tx)
.await?
}
_ => {
return Err(Error::InternalErr(format!(
"Matched unexpected JobKind ({:?}). This is a bug!",
kind
)))
}
};
job.runnable_id.replace(new_id.into());
}
job.runnable_id.replace(new_id.into());
}
// === RETRIEVE ACCUMULATED DEBOUNCE DATA ===
//
// For flows and apps, retrieve all nodes/components that were accumulated
// during the debounce window. This data comes from requests that were merged
// into this job instead of creating their own jobs.
//
// Scripts don't need this because they don't have nodes/components to relock.
if let Some(to_relock_field) = match &job.kind {
JobKind::FlowDependencies => Some("nodes_to_relock"),
JobKind::AppDependencies => Some("components_to_relock"),
_ => None, // Scripts don't use accumulated stale data
} {
tracing::debug!(
job_id = %job.id,
job_kind = ?job.kind,
field = %to_relock_field,
"Retrieving accumulated stale data from debounced requests"
);
if let Some(stale_data) = sqlx::query_scalar!(
"DELETE FROM debounce_stale_data WHERE job_id = $1 RETURNING to_relock",
&job.id
)
.fetch_optional(&mut *tx)
.await
.map_err(|e| {
tracing::error!(
error = %e,
job_id = %job.id,
"Failed to retrieve debounce_stale_data"
);
e
})?
.flatten()
{
// === RETRIEVE ACCUMULATED DEBOUNCE DATA ===
//
// For flows and apps, retrieve all nodes/components that were accumulated
// during the debounce window. This data comes from requests that were merged
// into this job instead of creating their own jobs.
//
// Scripts don't need this because they don't have nodes/components to relock.
if let Some(to_relock_field) = match &job.kind {
JobKind::FlowDependencies => Some("nodes_to_relock"),
JobKind::AppDependencies => Some("components_to_relock"),
_ => None, // Scripts don't use accumulated stale data
} {
tracing::debug!(
job_id = %job.id,
node_count = stale_data.len(),
nodes = ?stale_data,
"Retrieved accumulated nodes/components from {} debounced requests",
stale_data.len()
job_kind = ?job.kind,
field = %to_relock_field,
"Retrieving accumulated stale data from debounced requests"
);
// Replace the job's relock list with the accumulated data
// This ensures all nodes from all debounced requests are processed
if let Some(args) = job.args.as_mut() {
args.insert(to_relock_field.to_owned(), to_raw_value(&stale_data));
if let Some(stale_data) = sqlx::query_scalar!(
"DELETE FROM debounce_stale_data WHERE job_id = $1 RETURNING to_relock",
&job.id
)
.fetch_optional(&mut *tx)
.await
.map_err(|e| {
tracing::error!(
error = %e,
job_id = %job.id,
"Failed to retrieve debounce_stale_data"
);
e
})?
.flatten()
{
tracing::debug!(
field = %to_relock_field,
"Updated job args with accumulated debounce data"
job_id = %job.id,
node_count = stale_data.len(),
nodes = ?stale_data,
"Retrieved accumulated nodes/components from {} debounced requests",
stale_data.len()
);
// Replace the job's relock list with the accumulated data
// This ensures all nodes from all debounced requests are processed
if let Some(args) = job.args.as_mut() {
args.insert(to_relock_field.to_owned(), to_raw_value(&stale_data));
tracing::debug!(
field = %to_relock_field,
"Updated job args with accumulated debounce data"
);
}
} else {
tracing::trace!(
job_id = %job.id,
"No accumulated stale data found (no debounced requests or already cleaned up)"
);
}
} else {
tracing::trace!(
job_id = %job.id,
"No accumulated stale data found (no debounced requests or already cleaned up)"
);
}
}
// This will unblock pusher.
tx.commit().await?;
// This will unblock pusher.
tx.commit().await?;
Ok(())
}).await;
}
Ok(())
+1 -1
View File
@@ -1598,7 +1598,7 @@ pub async fn run_worker(
{
match timeout(
core::time::Duration::from_secs(10),
Box::pin(preprocess_dependency_job(pulled_job, &db)),
preprocess_dependency_job(pulled_job, &db),
)
.warn_after_seconds(2)
.await