diff --git a/backend/windmill-worker/src/js_eval.rs b/backend/windmill-worker/src/js_eval.rs index ac50b7c65c..b86327be28 100644 --- a/backend/windmill-worker/src/js_eval.rs +++ b/backend/windmill-worker/src/js_eval.rs @@ -340,6 +340,8 @@ lazy_static! { Regex::new(r#"(?m)(?Presults(?:(?:\.[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, + 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::>(); + 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 = vec![ @@ -679,7 +746,7 @@ pub async fn eval_fetch_timeout( Arc::new(BlobStore::default()), None, ), - deno_fetch::deno_fetch::init_ops::(deno_fetch_options), + deno_fetch::deno_fetch::init_ops::(fetch_options), deno_net::deno_net::init_ops::(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");