add proxy and useragent support for native scripts
This commit is contained in:
@@ -340,6 +340,8 @@ lazy_static! {
|
||||
Regex::new(r#"(?m)(?P<r>results(?:(?:\.[a-zA-Z_0-9]+)|(?:\[\".*?\"\])))"#).unwrap();
|
||||
static ref RE_FULL: Regex =
|
||||
Regex::new(r"(?m)^results\.([a-zA-Z_0-9]+)(?:\[(\d+)\])?((?:\.[a-zA-Z_0-9]+)+)?$").unwrap();
|
||||
static ref RE_PROXY: Regex =
|
||||
Regex::new(r"^(https?)://(([^:@\s]+):([^:@\s]+)@)?([^:@\s]+)(:(\d+))?$").unwrap();
|
||||
}
|
||||
|
||||
fn replace_with_await_result(expr: String) -> String {
|
||||
@@ -628,6 +630,53 @@ pub struct LogString {
|
||||
pub s: String,
|
||||
}
|
||||
|
||||
pub struct NativeAnnotation {
|
||||
pub useragent: Option<String>,
|
||||
pub proxy: Option<(String, Option<(String, String)>)>,
|
||||
}
|
||||
|
||||
pub fn get_annotation(inner_content: &str) -> NativeAnnotation {
|
||||
let mut res = NativeAnnotation { useragent: None, proxy: None };
|
||||
|
||||
#[cfg(feature = "enterprise")]
|
||||
let anns = inner_content
|
||||
.lines()
|
||||
.take_while(|x| x.starts_with("//"))
|
||||
.map(|x| x.to_string().trim_start_matches("//").trim().to_string())
|
||||
.collect_vec();
|
||||
|
||||
#[cfg(feature = "enterprise")]
|
||||
for ann in anns.iter() {
|
||||
if ann.starts_with("useragent") {
|
||||
res.useragent = Some(ann.trim_start_matches("useragent").trim().to_string());
|
||||
} else if ann.starts_with("proxy") {
|
||||
res.proxy = capture_proxy(ann.trim_start_matches("proxy").trim());
|
||||
}
|
||||
}
|
||||
res
|
||||
}
|
||||
|
||||
fn capture_proxy(s: &str) -> Option<(String, Option<(String, String)>)> {
|
||||
RE_PROXY.captures(s).map(|x| {
|
||||
(
|
||||
format!(
|
||||
"{}://{}{}",
|
||||
x.get(1).map(|x| x.as_str()).unwrap_or_default(),
|
||||
x.get(5).map(|x| x.as_str()).unwrap_or_default(),
|
||||
x.get(7)
|
||||
.map(|x| format!(":{}", x.as_str()))
|
||||
.unwrap_or_default(),
|
||||
),
|
||||
x.get(3).map(|y| {
|
||||
(
|
||||
y.as_str().to_string(),
|
||||
x.get(4).map(|x| x.as_str().to_string()).unwrap_or_default(),
|
||||
)
|
||||
}),
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn eval_fetch_timeout(
|
||||
env_code: String,
|
||||
ts_expr: String,
|
||||
@@ -653,22 +702,40 @@ pub async fn eval_fetch_timeout(
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let ann = get_annotation(&ts_expr);
|
||||
|
||||
let mut extra_logs = String::new();
|
||||
if ann.useragent.is_some() {
|
||||
extra_logs.push_str(&format!("useragent: {}\n", ann.useragent.as_ref().unwrap()));
|
||||
}
|
||||
if ann.proxy.is_some() {
|
||||
let (proxy, auth) = ann.proxy.as_ref().unwrap();
|
||||
extra_logs.push_str(&format!(
|
||||
"proxy: {proxy} (basic auth: {})\n",
|
||||
auth.is_some()
|
||||
));
|
||||
}
|
||||
|
||||
let result_f = tokio::task::spawn_blocking(move || {
|
||||
let ops = vec![op_get_static_args(), op_log()];
|
||||
let ext = Extension { name: "windmill", ops: ops.into(), ..Default::default() };
|
||||
|
||||
let deno_fetch_options = if let Some(cert_path) = env::var("DENO_CERT").ok() {
|
||||
let mut cert_store_provider = ContainerRootCertStoreProvider::new();
|
||||
cert_store_provider.add_certificate(cert_path)?;
|
||||
|
||||
deno_fetch::Options {
|
||||
root_cert_store_provider: Some(Arc::new(cert_store_provider)),
|
||||
user_agent: "windmill/beta".to_string(),
|
||||
proxy: None,
|
||||
..Default::default()
|
||||
}
|
||||
} else {
|
||||
Default::default()
|
||||
let fetch_options = deno_fetch::Options {
|
||||
root_cert_store_provider: if let Some(cert_path) = env::var("DENO_CERT").ok() {
|
||||
let mut cert_store_provider = ContainerRootCertStoreProvider::new();
|
||||
cert_store_provider.add_certificate(cert_path)?;
|
||||
Some(Arc::new(cert_store_provider))
|
||||
} else {
|
||||
None
|
||||
},
|
||||
user_agent: ann.useragent.unwrap_or_else(|| "windmill/beta".to_string()),
|
||||
proxy: ann.proxy.map(|x| deno_tls::Proxy {
|
||||
url: x.0,
|
||||
basic_auth: x
|
||||
.1
|
||||
.map(|(username, password)| deno_tls::BasicAuth { username, password }),
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let exts: Vec<Extension> = vec![
|
||||
@@ -679,7 +746,7 @@ pub async fn eval_fetch_timeout(
|
||||
Arc::new(BlobStore::default()),
|
||||
None,
|
||||
),
|
||||
deno_fetch::deno_fetch::init_ops::<PermissionsContainer>(deno_fetch_options),
|
||||
deno_fetch::deno_fetch::init_ops::<PermissionsContainer>(fetch_options),
|
||||
deno_net::deno_net::init_ops::<PermissionsContainer>(None, None),
|
||||
ext,
|
||||
];
|
||||
@@ -757,7 +824,7 @@ pub async fn eval_fetch_timeout(
|
||||
})
|
||||
});
|
||||
|
||||
let r = run_future_with_polling_update_job_poller(
|
||||
let (res, logs) = run_future_with_polling_update_job_poller(
|
||||
job_id,
|
||||
job_timeout,
|
||||
db,
|
||||
@@ -774,8 +841,8 @@ pub async fn eval_fetch_timeout(
|
||||
}
|
||||
e
|
||||
})?;
|
||||
*mem_peak = (r.0.get().len() / 1000) as i32;
|
||||
Ok(r)
|
||||
*mem_peak = (res.get().len() / 1000) as i32;
|
||||
Ok((res, format!("{extra_logs}{logs}")))
|
||||
}
|
||||
|
||||
const WINDMILL_CLIENT: &str = include_str!("./windmill-client.js");
|
||||
|
||||
Reference in New Issue
Block a user