* prometheus histogram for worker job timer hosts on :8001 * some new metrics in worker + field adds start_time_seconds, job_duration_seconds & jobs_failed * use tokio task_local to count job failures * METRICS_ADDR environment variable off by default true defaults to 0.0.0.0:8001 otherwise expects a socket address * pass metrics as args instead of task local
131 lines
4.2 KiB
Rust
131 lines
4.2 KiB
Rust
/*
|
|
* Author: Ruben Fiszel
|
|
* Copyright: Windmill Labs, Inc 2022
|
|
* This file and its contents are licensed under the AGPLv3 License.
|
|
* Please see the included NOTICE for copyright information and
|
|
* LICENSE-AGPL for a copy of the license.
|
|
*/
|
|
|
|
use std::net::SocketAddr;
|
|
|
|
use dotenv::dotenv;
|
|
|
|
#[tokio::main]
|
|
async fn main() -> anyhow::Result<()> {
|
|
dotenv().ok();
|
|
|
|
windmill::initialize_tracing().await?;
|
|
|
|
let db = windmill::connect_db().await?;
|
|
|
|
let num_workers = std::env::var("NUM_WORKERS")
|
|
.ok()
|
|
.and_then(|x| x.parse::<i32>().ok())
|
|
.unwrap_or(windmill::DEFAULT_NUM_WORKERS as i32);
|
|
|
|
let metrics_addr: Option<SocketAddr> = std::env::var("METRICS_ADDR")
|
|
.ok()
|
|
.map(|s| {
|
|
s.parse::<bool>()
|
|
.map(|b| b.then(|| SocketAddr::from(([0, 0, 0, 0], 8001))))
|
|
.or_else(|_| s.parse::<SocketAddr>().map(Some))
|
|
})
|
|
.transpose()?
|
|
.flatten();
|
|
|
|
let (server_mode, monitor_mode, migrate_db) = (true, true, true);
|
|
|
|
if migrate_db {
|
|
windmill::migrate_db(&db).await?;
|
|
}
|
|
|
|
let (tx, rx) = tokio::sync::broadcast::channel::<()>(3);
|
|
let shutdown_signal = windmill::shutdown_signal(tx.clone());
|
|
|
|
if server_mode || monitor_mode || num_workers > 0 {
|
|
let addr = SocketAddr::from(([0, 0, 0, 0], 8000));
|
|
|
|
let timeout = std::env::var("TIMEOUT")
|
|
.ok()
|
|
.and_then(|x| x.parse::<i32>().ok())
|
|
.unwrap_or(windmill::DEFAULT_TIMEOUT);
|
|
|
|
let server_f = async {
|
|
if server_mode {
|
|
windmill::run_server(
|
|
db.clone(),
|
|
addr,
|
|
&std::env::var("BASE_URL").unwrap_or("http://localhost".to_string()),
|
|
windmill::EmailSender {
|
|
from: "bot@windmill.dev".to_string(),
|
|
server: "smtp.gmail.com".to_string(),
|
|
password: std::env::var("SMTP_PASSWORD").unwrap_or("NOPASS".to_string()),
|
|
},
|
|
rx,
|
|
)
|
|
.await?;
|
|
}
|
|
Ok(()) as anyhow::Result<()>
|
|
};
|
|
|
|
let base_url = std::env::var("BASE_INTERNAL_URL")
|
|
.unwrap_or_else(|_| "http://missing-base-url".to_string());
|
|
|
|
let workers_f = async {
|
|
if num_workers > 0 {
|
|
let sleep_queue = std::env::var("SLEEP_QUEUE")
|
|
.ok()
|
|
.and_then(|x| x.parse::<u64>().ok())
|
|
.unwrap_or(windmill::DEFAULT_SLEEP_QUEUE);
|
|
let disable_nuser = std::env::var("DISABLE_NUSER")
|
|
.ok()
|
|
.and_then(|x| x.parse::<bool>().ok())
|
|
.unwrap_or(false);
|
|
let disable_nsjail = std::env::var("DISABLE_NSJAIL")
|
|
.ok()
|
|
.and_then(|x| x.parse::<bool>().ok())
|
|
.unwrap_or(false);
|
|
|
|
tracing::info!(
|
|
"DISABLE_NSJAIL: {disable_nsjail}, DISABLE_NUSER: {disable_nuser}, BASE_URL: \
|
|
{base_url}, SLEEP_QUEUE: {sleep_queue}, NUM_WORKERS: {num_workers}, TIMEOUT: \
|
|
{timeout}"
|
|
);
|
|
windmill::run_workers(
|
|
db.clone(),
|
|
addr,
|
|
timeout,
|
|
num_workers,
|
|
sleep_queue,
|
|
base_url,
|
|
disable_nuser,
|
|
disable_nsjail,
|
|
tx.clone(),
|
|
)
|
|
.await?;
|
|
}
|
|
Ok(()) as anyhow::Result<()>
|
|
};
|
|
|
|
let monitor_f = async {
|
|
if monitor_mode {
|
|
windmill::monitor_db(&db, timeout, tx.clone());
|
|
}
|
|
Ok(()) as anyhow::Result<()>
|
|
};
|
|
|
|
let metrics_f = async {
|
|
match metrics_addr {
|
|
Some(addr) => windmill::serve_metrics(addr, tx.subscribe())
|
|
.await
|
|
.map_err(anyhow::Error::from),
|
|
None => Ok(()),
|
|
}
|
|
};
|
|
|
|
futures::try_join!(shutdown_signal, server_f, workers_f, monitor_f, metrics_f)?;
|
|
}
|
|
|
|
Ok(())
|
|
}
|