diff --git a/Dockerfile b/Dockerfile index 9312c7e87c..70f8e2d852 100644 --- a/Dockerfile +++ b/Dockerfile @@ -59,13 +59,14 @@ FROM rust_base AS planner COPY ./openflow.openapi.yaml /openflow.openapi.yaml COPY ./backend ./ -RUN CARGO_NET_GIT_FETCH_WITH_CLI=true cargo chef prepare --recipe-path recipe.json +RUN CARGO_NET_GIT_FETCH_WITH_CLI=true cargo chef prepare --recipe-path recipe.json FROM rust_base AS builder +ARG features="" COPY --from=planner /windmill/recipe.json recipe.json -RUN CARGO_NET_GIT_FETCH_WITH_CLI=true cargo chef cook --release --recipe-path recipe.json +RUN CARGO_NET_GIT_FETCH_WITH_CLI=true cargo chef cook --release --features "$features" --recipe-path recipe.json COPY ./openflow.openapi.yaml /openflow.openapi.yaml COPY ./backend ./ @@ -73,7 +74,7 @@ COPY ./backend ./ COPY --from=frontend /frontend /frontend COPY .git/ .git/ -RUN CARGO_NET_GIT_FETCH_WITH_CLI=true cargo build --release +RUN CARGO_NET_GIT_FETCH_WITH_CLI=true cargo build --release --features "$features" FROM python:3.11.0-slim-buster @@ -81,9 +82,15 @@ FROM python:3.11.0-slim-buster ARG APP=/usr/src/app RUN apt-get update \ - && apt-get install -y ca-certificates wget curl git jq libprotobuf-dev libnl-route-3-dev build-essential \ + && apt-get install -y ca-certificates wget curl git jq libprotobuf-dev libnl-route-3-dev unzip \ + && apt-get install -y ca-certificates wget curl git jq libprotobuf-dev libnl-route-3-dev unzip build-essential \ && rm -rf /var/lib/apt/lists/* +RUN arch="$(dpkg --print-architecture)"; arch="${arch##*-}"; \ + curl -o rclone.zip "https://downloads.rclone.org/v1.60.1/rclone-v1.60.1-linux-$arch.zip"; \ + unzip -p rclone.zip rclone-v1.60.1-linux-$arch/rclone > /usr/bin/rclone; rm rclone.zip; \ + chown root:root /usr/bin/rclone; chmod 755 /usr/bin/rclone + RUN set -eux; \ arch="$(dpkg --print-architecture)"; arch="${arch##*-}"; \ url=; \ @@ -98,7 +105,7 @@ RUN set -eux; \ ;; \ *) echo >&2 "error: unsupported architecture '$arch' (likely packaging update needed)"; exit 1 ;; \ esac; \ - wget "https://golang.org/dl/$targz" && tar -C /usr/local -xzf "$targz" && rm "$targz"; + wget "https://golang.org/dl/$targz" -nv && tar -C /usr/local -xzf "$targz" && rm "$targz"; ENV PATH="${PATH}:/usr/local/go/bin" ENV GO_PATH=/usr/local/go/bin/go diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 56e80a4474..86355dae20 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -437,6 +437,27 @@ dependencies = [ "unicode-width", ] +[[package]] +name = "const_format" +version = "0.2.30" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7309d9b4d3d2c0641e018d449232f2e28f1b22933c137f157d3dbc14228b8c0e" +dependencies = [ + "const_format_proc_macros", + "konst", +] + +[[package]] +name = "const_format_proc_macros" +version = "0.2.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d897f47bf7270cf70d370f8f98c1abb6d2d4cf60a6845d30e05bfb90c6568650" +dependencies = [ + "proc-macro2", + "quote", + "unicode-xid", +] + [[package]] name = "convert_case" version = "0.4.0" @@ -1371,6 +1392,21 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "konst" +version = "0.2.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "330f0e13e6483b8c34885f7e6c9f19b1a7bd449c673fbb948a51c99d66ef74f4" +dependencies = [ + "konst_macro_rules", +] + +[[package]] +name = "konst_macro_rules" +version = "0.2.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4933f3f57a8e9d9da04db23fb153356ecaf00cbd14aee46279c33dc80925c37" + [[package]] name = "lalrpop" version = "0.19.8" @@ -4239,6 +4275,7 @@ dependencies = [ "anyhow", "async-recursion", "chrono", + "const_format", "deno_core", "dotenv", "futures", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index a9fc119080..dea616d4b9 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -27,6 +27,9 @@ edition = "2021" name = "windmill" path = "./src/main.rs" +[features] +enterprise = ["windmill-worker/enterprise"] + [dependencies] anyhow.workspace = true tokio.workspace = true @@ -134,3 +137,4 @@ futures = "^0" tokio-metrics = "0.1.0" lazy_static = "1.4.0" serde_derive = "1.0.147" +const_format = { version = "0.2", features = ["rust_1_64", "rust_1_51"] } diff --git a/backend/src/main.rs b/backend/src/main.rs index 9a31d656c7..8872fcc5d4 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -87,6 +87,10 @@ async fn main() -> anyhow::Result<()> { .ok() .and_then(|x| x.parse::().ok()) .unwrap_or(false); + let sync_bucket = std::env::var("S3_CACHE_BUCKET") + .ok() + .map(|e| Some(e)) + .unwrap_or(None); tracing::info!( "DISABLE_NSJAIL: {disable_nsjail}, DISABLE_NUSER: {disable_nuser}, BASE_URL: \ @@ -107,6 +111,7 @@ async fn main() -> anyhow::Result<()> { keep_job_dir, }, rx.resubscribe(), + sync_bucket, ) .await?; } @@ -160,6 +165,7 @@ pub async fn run_workers( sleep_queue: u64, worker_config: WorkerConfig, rx: tokio::sync::broadcast::Receiver<()>, + mut periodic_script: Option, ) -> anyhow::Result<()> { let instance_name = rd_string(5); let monitor = tokio_metrics::TaskMonitor::new(); @@ -180,6 +186,7 @@ pub async fn run_workers( let ip = ip.clone(); let rx = rx.resubscribe(); let worker_config = worker_config.clone(); + let wp = periodic_script.take(); handles.push(tokio::spawn(monitor.instrument(async move { tracing::info!(addr = %addr.to_string(), worker = %worker_name, "starting worker"); windmill_worker::run_worker( @@ -192,6 +199,7 @@ pub async fn run_workers( &ip, sleep_queue, worker_config, + wp, rx, ) .await diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 5f284c88d1..bdd8b22994 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -932,6 +932,7 @@ fn spawn_test_worker( ip, sleep_queue, worker_config, + None, rx, ) .await diff --git a/backend/windmill-worker/Cargo.toml b/backend/windmill-worker/Cargo.toml index c883ac7e6e..9774434c94 100644 --- a/backend/windmill-worker/Cargo.toml +++ b/backend/windmill-worker/Cargo.toml @@ -12,6 +12,7 @@ path = "./src/main.rs" [features] default = [] deno-lock = [] +enterprise = [] [dependencies] windmill-queue.workspace = true @@ -44,4 +45,5 @@ lazy_static.workspace = true chrono.workspace = true dotenv.workspace = true rand.workspace = true # TODO: Remove. only used by token creation hack. -deno_core.workspace = true \ No newline at end of file +deno_core.workspace = true +const_format.workspace = true diff --git a/backend/windmill-worker/src/main.rs b/backend/windmill-worker/src/main.rs index 81bf2bc5ef..f21761ff34 100644 --- a/backend/windmill-worker/src/main.rs +++ b/backend/windmill-worker/src/main.rs @@ -81,6 +81,10 @@ async fn main() -> anyhow::Result<()> { .ok() .and_then(|x| x.parse::().ok()) .unwrap_or(false); + let sync_bucket = std::env::var("S3_CACHE_BUCKET") + .ok() + .map(|e| Some(e)) + .unwrap_or(None); tracing::info!( "DISABLE_NSJAIL: {disable_nsjail}, DISABLE_NUSER: {disable_nuser}, BASE_URL: \ @@ -112,6 +116,7 @@ async fn main() -> anyhow::Result<()> { base_url, keep_job_dir, }, + sync_bucket, rx.resubscribe(), ) .await; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 5dce1e7be6..a241788e80 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -6,6 +6,7 @@ * LICENSE-AGPL for a copy of the license. */ +use const_format::concatcp; use itertools::Itertools; use sqlx::{Pool, Postgres, Transaction}; use std::{borrow::Borrow, collections::HashMap, io, panic, process::Stdio, time::Duration}; @@ -35,7 +36,7 @@ use tokio::{ use futures::{ future::{self, ready, FutureExt}, - stream::{self, StreamExt}, + stream, StreamExt, }; use async_recursion::async_recursion; @@ -47,6 +48,51 @@ use crate::{ }, }; +#[cfg(feature = "enterprise")] +async fn run_periodic_jobs(bucket: &str) { + tracing::info!("Running periodic jobs"); + + match Command::new("rclone") + .arg("bisync") + .arg(format!(":s3,env_auth=true:{}", bucket)) + .arg(ROOT_CACHE_DIR) + .arg("--size-only") + .arg("--fast-list") + .arg("--force") + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .spawn() + { + Ok(mut h) => { + h.wait().await.unwrap(); + } + Err(e) => tracing::warn!("Failed to run periodic job. Error: {:?}", e), + } +} + +#[cfg(feature = "enterprise")] +async fn run_periodic_jobs_initial(bucket: &str) { + tracing::info!("Running periodic jobs"); + + match Command::new("rclone") + .arg("bisync") + .arg(format!(":s3,env_auth=true:{}", bucket)) + .arg(ROOT_CACHE_DIR) + .arg("--size-only") + .arg("--fast-list") + .arg("--force") + .arg("--resync") + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .spawn() + { + Ok(mut h) => { + h.wait().await.unwrap(); + } + Err(e) => tracing::warn!("Failed to run periodic job. Error: {:?}", e), + } +} + #[tracing::instrument(level = "trace", skip_all)] pub async fn create_token_for_owner<'c>( mut tx: Transaction<'c, Postgres>, @@ -110,10 +156,12 @@ pub async fn create_token_for_owner<'c>( const TMP_DIR: &str = "/tmp/windmill"; const PIP_SUPERCACHE_DIR: &str = "/tmp/windmill/cache/pip_permanent"; -const PIP_CACHE_DIR: &str = "/tmp/windmill/cache/pip"; -const DENO_CACHE_DIR: &str = "/tmp/windmill/cache/deno"; -const GO_CACHE_DIR: &str = "/tmp/windmill/cache/go"; +const ROOT_CACHE_DIR: &str = "/tmp/windmill/cache/"; +const PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "pip"); +const DENO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "deno"); +const GO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "go"); const NUM_SECS_ENV_CHECK: u64 = 15; +const NUM_SECS_SYNC: u64 = 60 * 10; const DEFAULT_HEAVY_DEPS: [&str; 18] = [ "numpy", "pandas", @@ -192,6 +240,7 @@ pub async fn run_worker( ip: &str, sleep_queue: u64, worker_config: WorkerConfig, + sync_bucket: Option, mut rx: tokio::sync::broadcast::Receiver<()>, ) { let start_time = Instant::now(); @@ -220,7 +269,10 @@ pub async fn run_worker( ) .await; - let mut last_ping = Instant::now() - Duration::from_secs(NUM_SECS_ENV_CHECK + 1); + #[cfg(feature = "enterprise")] + let mut last_sync = Instant::now() - Duration::from_secs(NUM_SECS_ENV_CHECK + 1); + + let mut last_ping = Instant::now() - Duration::from_secs(NUM_SECS_SYNC + 1); insert_initial_ping(worker_instance, &worker_name, ip, db).await; @@ -287,6 +339,11 @@ pub async fn run_worker( }; WORKER_STARTED.inc(); + #[cfg(feature = "enterprise")] + if let Some(ref s) = sync_bucket.clone() { + run_periodic_jobs_initial(&s).await; + } + let (same_worker_tx, mut same_worker_rx) = mpsc::channel::(5); loop { @@ -309,6 +366,14 @@ pub async fn run_worker( last_ping = Instant::now(); } + #[cfg(feature = "enterprise")] + if last_sync.elapsed().as_secs() > NUM_SECS_SYNC { + if let Some(ref s) = sync_bucket.clone() { + run_periodic_jobs(&s).await; + } + last_sync = Instant::now(); + } + let (do_break, next_job) = async { tokio::select! { biased;