mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-11 08:07:15 +00:00
hub script fetch retry (#5379)
* hub script fetch retry * use backon * oups * Update backend/windmill-common/src/scripts.rs Co-authored-by: ellipsis-dev[bot] <65095814+ellipsis-dev[bot]@users.noreply.github.com> * retry whole logic * nits --------- Co-authored-by: ellipsis-dev[bot] <65095814+ellipsis-dev[bot]@users.noreply.github.com>
This commit is contained in:
co-authored by
ellipsis-dev[bot]
parent
7bf9e25ede
commit
d30979d04e
Generated
+1
@@ -13682,6 +13682,7 @@ dependencies = [
|
||||
"aws-config",
|
||||
"aws-sdk-sts",
|
||||
"axum",
|
||||
"backon",
|
||||
"bytes",
|
||||
"chrono",
|
||||
"chrono-tz 0.10.1",
|
||||
|
||||
@@ -61,6 +61,7 @@ async-stream.workspace = true
|
||||
const_format.workspace = true
|
||||
crc.workspace = true
|
||||
windmill-macros.workspace = true
|
||||
backon.workspace = true
|
||||
|
||||
semver.workspace = true
|
||||
croner = "2.0.6"
|
||||
|
||||
@@ -19,6 +19,8 @@ use crate::{
|
||||
|
||||
use crate::worker::HUB_CACHE_DIR;
|
||||
use anyhow::Context;
|
||||
use backon::ConstantBuilder;
|
||||
use backon::{BackoffBuilder, Retryable};
|
||||
use serde::de::Error as _;
|
||||
use serde::{ser::SerializeSeq, Deserialize, Deserializer, Serialize};
|
||||
|
||||
@@ -415,6 +417,8 @@ pub async fn get_hub_script_by_path(
|
||||
Some(db),
|
||||
)
|
||||
.await?
|
||||
.error_for_status()
|
||||
.map_err(to_anyhow)?
|
||||
.text()
|
||||
.await
|
||||
.map_err(to_anyhow);
|
||||
@@ -440,6 +444,8 @@ pub async fn get_hub_script_by_path(
|
||||
Some(db),
|
||||
)
|
||||
.await?
|
||||
.error_for_status()
|
||||
.map_err(to_anyhow)?
|
||||
.text()
|
||||
.await
|
||||
.map_err(to_anyhow)?;
|
||||
@@ -494,49 +500,67 @@ async fn get_full_hub_script_by_path_inner(
|
||||
) -> crate::error::Result<HubScript> {
|
||||
let hub_base_url = HUB_BASE_URL.read().await.clone();
|
||||
|
||||
let result = http_get_from_hub(
|
||||
http_client,
|
||||
&format!("{}/raw2/{}", hub_base_url, path),
|
||||
true,
|
||||
None,
|
||||
db,
|
||||
)
|
||||
.await?
|
||||
.json::<HubScript>()
|
||||
.await
|
||||
.context("Decoding hub response to script");
|
||||
let response = (|| async {
|
||||
let response = http_get_from_hub(
|
||||
http_client,
|
||||
&format!("{}/raw2/{}", hub_base_url, path),
|
||||
true,
|
||||
None,
|
||||
db,
|
||||
)
|
||||
.await
|
||||
.and_then(|r| r.error_for_status().map_err(|e| to_anyhow(e).into()));
|
||||
|
||||
match result {
|
||||
Ok(result) => Ok(result),
|
||||
Err(e) => {
|
||||
if hub_base_url != DEFAULT_HUB_BASE_URL
|
||||
&& path
|
||||
.split("/")
|
||||
.next()
|
||||
.is_some_and(|x| x.parse::<i32>().is_ok_and(|x| x < 10_000_000))
|
||||
{
|
||||
tracing::info!(
|
||||
"Not found on private hub, fallback to default hub for {}",
|
||||
path
|
||||
);
|
||||
let value = http_get_from_hub(
|
||||
http_client,
|
||||
&format!("{}/raw2/{}", DEFAULT_HUB_BASE_URL, path),
|
||||
true,
|
||||
None,
|
||||
db,
|
||||
)
|
||||
.await?
|
||||
.json::<HubScript>()
|
||||
.await
|
||||
.context("Decoding hub response to script")?;
|
||||
|
||||
Ok(value)
|
||||
} else {
|
||||
Err(e)?
|
||||
match response {
|
||||
Ok(response) => Ok(response),
|
||||
Err(e) => {
|
||||
if hub_base_url != DEFAULT_HUB_BASE_URL
|
||||
&& path
|
||||
.split("/")
|
||||
.next()
|
||||
.is_some_and(|x| x.parse::<i32>().is_ok_and(|x| x < 10_000_000))
|
||||
{
|
||||
// TODO: should only fallback to default hub if status is 404 (hub returns 500 currently)
|
||||
tracing::info!(
|
||||
"Not found on private hub, fallback to default hub for {}",
|
||||
path
|
||||
);
|
||||
http_get_from_hub(
|
||||
http_client,
|
||||
&format!("{}/raw2/{}", DEFAULT_HUB_BASE_URL, path),
|
||||
true,
|
||||
None,
|
||||
db,
|
||||
)
|
||||
.await?
|
||||
.error_for_status()
|
||||
.map_err(|e| to_anyhow(e).into())
|
||||
} else {
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
.retry(
|
||||
ConstantBuilder::default()
|
||||
.with_delay(std::time::Duration::from_secs(5))
|
||||
.with_max_times(2)
|
||||
.build(),
|
||||
)
|
||||
.notify(|err, dur| {
|
||||
tracing::warn!(
|
||||
"Could not get hub script at path {path}, retrying in {dur:#?}, err: {err:#?}"
|
||||
);
|
||||
})
|
||||
.sleep(tokio::time::sleep)
|
||||
.await?;
|
||||
|
||||
let script = response
|
||||
.json::<HubScript>()
|
||||
.await
|
||||
.context(format!("Decoding hub response for script at path {path}"))?;
|
||||
|
||||
Ok(script)
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Serialize)]
|
||||
|
||||
Reference in New Issue
Block a user