mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-05 00:02:24 +00:00
* feat: let a workspace withdraw operator schedule and trigger writes Operators can create, edit and delete schedules and triggers today through the API, CLI and MCP, while the operator_settings flags beside them only hide those pages. An admin who wants operators to see what is scheduled without letting them change it cannot express that. Add manage_schedules and manage_triggers as enforced settings, gated at the schedule handlers and at the generic TriggerCrud routes so every trigger kind is covered by one check. They name capabilities operators already hold, so they are granted unless withdrawn, and absence has to mean "never configured" rather than a value. The read coalesces to true; the update endpoint merges into the stored jsonb with the two fields as Option<bool>, so an omitted key keeps what is stored. operator_settings is git-synced as a whole object, so a settings file written before these keys existed reaches the endpoint on every pull, and a serde or SQL default of either polarity would turn that pull into a silent withdrawal or restoration. The rights are read through a per-process cache, so withdrawing one publishes a notify_event that drops the entry on every replica. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Dsf6VC4MVLisiEoeQkgbr4 * feat: enforce operator write rights on the router and in the UI Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: close the capture gap and gate the trigger editors' write actions Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: gate acl writes and the native trigger drawer behind manage rights Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: refuse operator writes with 403 and gate sharing at the drawer Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * perf: resolve identity in the operator write gate only for writes Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: gate the suspended-jobs actions and stop the route check refusing reads Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: explain the empty-state create button when operator writes are withdrawn Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: audit operator settings changes and fold path writes into native rows Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: open locked editors read-only and group the operator settings Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: skip email and azure lookups on editor open while triggers are locked Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * docs: state each operator-rights rationale once in comments Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: address CI review findings on operator write rights Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: keep capture move gated and skip it in the builders while locked Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: keep admin and operator exclusive when setting a workspace role Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * refactor: use the shared section component for operator settings groups Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com> Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
1487 lines
57 KiB
Rust
1487 lines
57 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 crate::db::ApiAuthed;
|
|
#[cfg(feature = "embedding")]
|
|
use crate::embeddings::load_embeddings_db;
|
|
#[cfg(feature = "oauth2")]
|
|
use crate::oauth2_oss::SlackVerifier;
|
|
#[cfg(feature = "smtp")]
|
|
use crate::smtp_server_oss::SmtpServer;
|
|
use axum::middleware::from_fn_with_state;
|
|
#[cfg(feature = "enterprise")]
|
|
use windmill_api_auth::ee_oss::ExternalJwks;
|
|
use windmill_api_auth::gate_operator_writes;
|
|
use windmill_common::workspaces::ManageKind;
|
|
use windmill_store::resources::public_service;
|
|
|
|
#[cfg(feature = "mcp")]
|
|
use crate::mcp::{extract_and_store_workspace_id, setup_mcp_server};
|
|
use crate::triggers::start_all_listeners;
|
|
use tower_http::catch_panic::CatchPanicLayer;
|
|
|
|
use crate::tracing_init::MyOnFailure;
|
|
use crate::{
|
|
s3_log_batching::{s3_proxy_log_middleware, FLUSH_INTERVAL_MS},
|
|
tracing_init::{MyMakeSpan, MyOnResponse},
|
|
users::OptAuthed,
|
|
webhook_util::WebhookShared,
|
|
};
|
|
#[cfg(feature = "agent_worker_server")]
|
|
use windmill_api_agent_workers::AgentCache;
|
|
|
|
use anyhow::Context;
|
|
use argon2::Argon2;
|
|
use axum::body::Body;
|
|
use axum::extract::DefaultBodyLimit;
|
|
use axum::http::HeaderValue;
|
|
use axum::response::Response;
|
|
use axum::serve::ListenerExt;
|
|
use axum::{middleware::from_extractor, routing::get, routing::post, Extension, Json, Router};
|
|
use db::DB;
|
|
use tokio::task::JoinHandle;
|
|
use windmill_common::global_settings::load_value_from_global_settings;
|
|
use windmill_common::global_settings::EMAIL_DOMAIN_SETTING;
|
|
use windmill_common::worker::HUB_CACHE_DIR;
|
|
|
|
use std::fs::DirBuilder;
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
use tokio::sync::RwLock;
|
|
use tower::ServiceBuilder;
|
|
use tower_cookies::CookieManagerLayer;
|
|
use tower_http::{
|
|
cors::{Any, CorsLayer},
|
|
trace::TraceLayer,
|
|
};
|
|
use windmill_common::db::UserDB;
|
|
use windmill_common::worker::CLOUD_HOSTED;
|
|
#[allow(unused_imports)]
|
|
pub(crate) use windmill_common::BASE_URL;
|
|
use windmill_common::{utils::GIT_VERSION, INSTANCE_NAME};
|
|
|
|
use crate::scim_oss::has_scim_token;
|
|
use windmill_common::error::AppError;
|
|
|
|
mod ai;
|
|
#[cfg(feature = "private")]
|
|
mod ai_free_tier_ee;
|
|
mod ai_free_tier_oss;
|
|
#[cfg(feature = "parquet")]
|
|
mod ai_sessions;
|
|
#[cfg(feature = "parquet")]
|
|
pub use ai_sessions::sweep_expired_ai_session_backups;
|
|
mod ai_shared_artifacts;
|
|
mod apps;
|
|
mod apps_raw_bundle;
|
|
pub use apps::invalidate_app_policy_cache;
|
|
pub mod args;
|
|
mod audit;
|
|
pub mod auth;
|
|
#[cfg(all(feature = "private", feature = "parquet"))]
|
|
pub mod azure_proxy_ee;
|
|
mod azure_proxy_oss;
|
|
mod capture;
|
|
mod concurrency_groups;
|
|
mod csrf;
|
|
mod db;
|
|
mod db_health;
|
|
mod dbt;
|
|
mod docs;
|
|
mod drafts;
|
|
|
|
#[cfg(feature = "private")]
|
|
pub mod ee;
|
|
pub mod ee_oss;
|
|
pub mod embeddings;
|
|
mod favorite;
|
|
pub mod flows;
|
|
mod folder_history;
|
|
mod folders;
|
|
mod granular_acls;
|
|
mod group_history;
|
|
mod groups;
|
|
mod health;
|
|
mod hub_publish;
|
|
#[cfg(feature = "private")]
|
|
pub mod indexer_ee;
|
|
mod indexer_oss;
|
|
mod integration;
|
|
mod internal_db;
|
|
mod live_migrations;
|
|
mod runnables;
|
|
#[cfg(all(feature = "private", feature = "parquet"))]
|
|
pub mod s3_proxy_ee;
|
|
mod s3_proxy_oss;
|
|
#[cfg(all(feature = "private", feature = "parquet"))]
|
|
pub mod storage_list_ee;
|
|
#[cfg(feature = "parquet")]
|
|
mod storage_list_oss;
|
|
mod workspace_dependencies;
|
|
|
|
mod ai_evals;
|
|
mod approvals;
|
|
#[cfg(all(feature = "enterprise", feature = "private"))]
|
|
pub mod apps_ee;
|
|
#[cfg(feature = "enterprise")]
|
|
mod apps_oss;
|
|
#[cfg(all(feature = "enterprise", feature = "private"))]
|
|
pub mod git_sync_ee;
|
|
#[cfg(feature = "enterprise")]
|
|
mod git_sync_oss;
|
|
#[cfg(all(feature = "parquet", feature = "private"))]
|
|
pub mod job_helpers_ee;
|
|
#[cfg(feature = "parquet")]
|
|
mod job_helpers_oss;
|
|
pub mod job_metrics;
|
|
pub mod jobs;
|
|
pub mod jobs_export;
|
|
#[cfg(all(feature = "oauth2", feature = "private"))]
|
|
pub mod oauth2_ee;
|
|
#[cfg(feature = "oauth2")]
|
|
pub mod oauth2_oss;
|
|
#[cfg(feature = "private")]
|
|
pub mod oidc_ee;
|
|
mod oidc_oss;
|
|
mod path_autocomplete;
|
|
mod raw_apps;
|
|
mod resources;
|
|
#[cfg(feature = "private")]
|
|
pub mod saml_ee;
|
|
mod saml_oss;
|
|
#[cfg(feature = "private")]
|
|
pub mod scim_ee;
|
|
mod scim_oss;
|
|
mod scripts;
|
|
mod secret_backend_ext;
|
|
mod service_logs;
|
|
mod slack_approvals;
|
|
#[cfg(all(feature = "smtp", feature = "private"))]
|
|
pub mod smtp_server_ee;
|
|
#[cfg(feature = "smtp")]
|
|
mod smtp_server_oss;
|
|
#[cfg(feature = "private")]
|
|
pub mod teams_approvals_ee;
|
|
mod teams_approvals_oss;
|
|
mod workspace_shared_ui;
|
|
|
|
#[cfg(feature = "native_trigger")]
|
|
pub mod native_triggers;
|
|
mod offboarding;
|
|
mod public_app_layer;
|
|
mod public_app_rate_limit;
|
|
mod s3_log_batching;
|
|
mod static_assets;
|
|
#[cfg(all(feature = "stripe", feature = "enterprise", feature = "private"))]
|
|
pub mod stripe_ee;
|
|
#[cfg(all(feature = "stripe", feature = "enterprise"))]
|
|
mod stripe_oss;
|
|
#[cfg(feature = "private")]
|
|
pub mod teams_cache_ee;
|
|
mod teams_cache_oss;
|
|
#[cfg(feature = "private")]
|
|
pub mod teams_ee;
|
|
mod teams_oss;
|
|
mod token;
|
|
mod tracing_init;
|
|
mod trash;
|
|
mod trigger_history;
|
|
pub mod triggers;
|
|
mod users;
|
|
#[cfg(feature = "private")]
|
|
pub mod users_ee;
|
|
mod users_oss;
|
|
mod utils;
|
|
mod variables;
|
|
#[cfg(feature = "private")]
|
|
pub mod volumes_ee;
|
|
mod volumes_oss;
|
|
pub mod webhook_util;
|
|
mod workspaces;
|
|
#[cfg(feature = "private")]
|
|
pub mod workspaces_ee;
|
|
mod workspaces_export;
|
|
|
|
#[cfg(feature = "mcp")]
|
|
mod mcp_tools;
|
|
|
|
#[cfg(feature = "mcp")]
|
|
mod mcp;
|
|
#[cfg(all(feature = "mcp", feature = "private"))]
|
|
mod mcp_oauth_ee;
|
|
#[cfg(feature = "mcp")]
|
|
mod mcp_oauth_oss;
|
|
|
|
pub use apps::EditApp;
|
|
pub const DEFAULT_BODY_LIMIT: usize = 2097152 * 100; // 200MB
|
|
|
|
lazy_static::lazy_static! {
|
|
|
|
pub static ref REQUEST_SIZE_LIMIT: Arc<RwLock<usize>> = Arc::new(RwLock::new(DEFAULT_BODY_LIMIT));
|
|
|
|
pub static ref SCIM_TOKEN: Arc<RwLock<Option<String>>> = Arc::new(RwLock::new(None));
|
|
pub static ref SAML_METADATA: Arc<RwLock<Option<String>>> = Arc::new(RwLock::new(None));
|
|
|
|
|
|
// COOKIE_DOMAIN and IS_SECURE are now in windmill_common::utils
|
|
}
|
|
|
|
pub use windmill_common::utils::HTTP_CLIENT_PERMISSIVE as HTTP_CLIENT;
|
|
|
|
pub use windmill_common::utils::{COOKIE_DOMAIN, IS_SECURE};
|
|
|
|
pub use windmill_api_debug::reload_debug_signing_key;
|
|
|
|
#[cfg(feature = "oauth2")]
|
|
pub use windmill_oauth::OAUTH_CLIENTS;
|
|
|
|
#[cfg(feature = "oauth2")]
|
|
lazy_static::lazy_static! {
|
|
pub static ref SLACK_SIGNING_SECRET: Option<SlackVerifier> = std::env::var("SLACK_SIGNING_SECRET")
|
|
.ok()
|
|
.map(|x| SlackVerifier::new(x).unwrap());
|
|
|
|
// Comma-separated Slack v2 OAuth bot scopes requested when connecting a workspace.
|
|
// Must be a subset of the bot scopes declared in the Slack app manifest. Default matches
|
|
// Windmill's recommended manifest at docs.windmill.dev/docs/misc/setup_oauth.
|
|
pub static ref SLACK_OAUTH_SCOPES: String = std::env::var("SLACK_OAUTH_SCOPES")
|
|
.unwrap_or_else(|_| "commands,chat:write,chat:write.public,channels:join,files:write,app_mentions:read,im:history,im:read".to_string());
|
|
}
|
|
|
|
// Compliance with cloud events spec.
|
|
pub async fn add_webhook_allowed_origin(
|
|
req: axum::extract::Request,
|
|
next: axum::middleware::Next,
|
|
) -> axum::response::Response {
|
|
if req.method() == http::Method::OPTIONS {
|
|
if let Some(webhook_request_origin) = req.headers().get("Webhook-Request-Origin") {
|
|
let webhook_request_origin = webhook_request_origin.clone();
|
|
let mut response = next.run(req).await;
|
|
|
|
response
|
|
.headers_mut()
|
|
.insert("Webhook-Allowed-Origin", webhook_request_origin);
|
|
response
|
|
.headers_mut()
|
|
.insert("Webhook-Allowed-Rate", HeaderValue::from_static("*"));
|
|
return response;
|
|
}
|
|
}
|
|
next.run(req).await
|
|
}
|
|
|
|
/// Scope the request in the deploy origin it declares, so the fork tally can tell
|
|
/// a write applied by a sync from one authored in the workspace. Entered for
|
|
/// every request, including unmarked ones: `deploy_origin::current` reads the
|
|
/// scope's presence as "the caller is serving this write", which is what
|
|
/// separates a request from a worker that tallies someone else's.
|
|
async fn set_deploy_origin(
|
|
req: axum::extract::Request,
|
|
next: axum::middleware::Next,
|
|
) -> axum::response::Response {
|
|
let origin = req
|
|
.headers()
|
|
.get(windmill_common::deploy_origin::DEPLOY_ORIGIN_HEADER)
|
|
.and_then(|v| v.to_str().ok())
|
|
.map(windmill_common::deploy_origin::DeployOrigin::from_header_value)
|
|
.unwrap_or(windmill_common::deploy_origin::DeployOrigin::Authored);
|
|
windmill_common::deploy_origin::scope(origin, next.run(req)).await
|
|
}
|
|
|
|
/// Scope the request in the client kind it declares, so a trigger mutation can
|
|
/// be attributed to the CLI rather than to a bare API call. Entered for every
|
|
/// request, undeclared ones included: `TriggerSource::of_request` reads the
|
|
/// scope's absence as "no request is being served", which is what separates a
|
|
/// caller from a worker disabling a trigger on its own.
|
|
async fn set_request_client(
|
|
req: axum::extract::Request,
|
|
next: axum::middleware::Next,
|
|
) -> axum::response::Response {
|
|
let client = req
|
|
.headers()
|
|
.get(windmill_common::trigger_history::CLIENT_HEADER)
|
|
.and_then(|v| v.to_str().ok())
|
|
.and_then(windmill_common::trigger_history::client_from_header);
|
|
windmill_common::trigger_history::scope_client(client, next.run(req)).await
|
|
}
|
|
|
|
#[cfg(not(feature = "tantivy"))]
|
|
type IndexReader = ();
|
|
|
|
#[cfg(not(feature = "tantivy"))]
|
|
type ServiceLogIndexReader = ();
|
|
|
|
#[cfg(feature = "tantivy")]
|
|
type IndexReader = windmill_indexer::completed_runs_oss::IndexReader;
|
|
#[cfg(feature = "tantivy")]
|
|
type ServiceLogIndexReader = windmill_indexer::service_logs_oss::ServiceLogIndexReader;
|
|
|
|
/// Worker name derived from the agent JWT token, used to authenticate volume operations.
|
|
/// Defined unconditionally so volume endpoint handlers can reference it regardless of
|
|
/// whether agent_worker_server is enabled (the extension is only populated on the agent path).
|
|
#[derive(Clone)]
|
|
pub struct AgentWorkerName(pub String);
|
|
|
|
/// Middleware that injects a synthetic `ApiAuthed` and JWT-derived worker name
|
|
/// into request extensions.
|
|
///
|
|
/// Used for volume proxy endpoints under the agent_workers path, where the
|
|
/// agent JWT auth layer has already validated the request. The volume handlers
|
|
/// need `ApiAuthed` to resolve the workspace S3 client, but the agent JWT
|
|
/// format is incompatible with the standard auth extractor.
|
|
///
|
|
/// The worker name is extracted from the JWT claims rather than trusting
|
|
/// self-reported values in request bodies/query params.
|
|
#[cfg(feature = "agent_worker_server")]
|
|
async fn inject_agent_authed(
|
|
request: axum::extract::Request,
|
|
next: axum::middleware::Next,
|
|
) -> Response {
|
|
let mut request = request;
|
|
|
|
// Extract worker name from agent JWT via AgentCache
|
|
// (OSS returns None; EE decodes the JWT and returns the worker name)
|
|
{
|
|
let extracted = {
|
|
let token = request
|
|
.headers()
|
|
.get(axum::http::header::AUTHORIZATION)
|
|
.and_then(|v| v.to_str().ok())
|
|
.and_then(|s| s.strip_prefix("Bearer ").map(|t| t.to_string()));
|
|
let cache = request.extensions().get::<Arc<AgentCache>>().cloned();
|
|
let db = request.extensions().get::<DB>().cloned();
|
|
match (token, cache, db) {
|
|
(Some(token), Some(cache), Some(db)) => Some((token, cache, db)),
|
|
_ => None,
|
|
}
|
|
};
|
|
|
|
if let Some((token, cache, db)) = extracted {
|
|
if let Some(worker_name) = cache.extract_worker_name(&token, &db).await {
|
|
request
|
|
.extensions_mut()
|
|
.insert(AgentWorkerName(worker_name));
|
|
}
|
|
}
|
|
}
|
|
|
|
request
|
|
.extensions_mut()
|
|
.insert(windmill_api_auth::OptJobAuthed {
|
|
authed: ApiAuthed {
|
|
email: "agent-worker@windmill.dev".to_string(),
|
|
username: "agent-worker".to_string(),
|
|
is_admin: true,
|
|
is_operator: false,
|
|
groups: Vec::new(),
|
|
folders: Vec::new(),
|
|
scopes: None,
|
|
username_override: None,
|
|
username_override_is_token_label: false,
|
|
is_session_token: false,
|
|
token_prefix: None,
|
|
read_only: false,
|
|
job_id: None,
|
|
credential_expiry: None,
|
|
},
|
|
job_id: None,
|
|
});
|
|
next.run(request).await
|
|
}
|
|
|
|
pub async fn run_server(
|
|
db: DB,
|
|
job_index_reader: Option<IndexReader>,
|
|
log_index_reader: Option<ServiceLogIndexReader>,
|
|
listener: tokio::net::TcpListener,
|
|
mut killpill_rx: tokio::sync::broadcast::Receiver<()>,
|
|
port_tx: tokio::sync::oneshot::Sender<String>,
|
|
server_mode: bool,
|
|
mcp_mode: bool,
|
|
_base_internal_url: String,
|
|
name: Option<String>,
|
|
) -> anyhow::Result<()> {
|
|
let user_db = UserDB::new(db.clone());
|
|
|
|
for x in [&*HUB_CACHE_DIR] {
|
|
DirBuilder::new()
|
|
.recursive(true)
|
|
.create(x)
|
|
.expect("could not create initial server dir");
|
|
}
|
|
|
|
#[cfg(feature = "enterprise")]
|
|
let ext_jwks = ExternalJwks::load().await;
|
|
let auth_cache = Arc::new(crate::auth::AuthCache::new(
|
|
db.clone(),
|
|
std::env::var("SUPERADMIN_SECRET").ok(),
|
|
#[cfg(feature = "enterprise")]
|
|
ext_jwks,
|
|
));
|
|
let argon2 = Arc::new(Argon2::default());
|
|
|
|
// Initialize debug signing key for debugger authentication
|
|
windmill_api_debug::init_debug_signing_key().await;
|
|
|
|
let disable_response_logs = std::env::var("DISABLE_RESPONSE_LOGS")
|
|
.ok()
|
|
.map(|x| x == "true")
|
|
.unwrap_or(false);
|
|
|
|
let middleware_stack = ServiceBuilder::new()
|
|
.layer(Extension(db.clone()))
|
|
.layer(Extension(user_db.clone()))
|
|
.layer(Extension(auth_cache.clone()))
|
|
.layer(Extension(job_index_reader))
|
|
.layer(Extension(log_index_reader))
|
|
// .layer(Extension(index_writer))
|
|
.layer(CookieManagerLayer::new())
|
|
.layer(Extension(WebhookShared::new(
|
|
killpill_rx.resubscribe(),
|
|
db.clone(),
|
|
)))
|
|
.layer(DefaultBodyLimit::max(
|
|
REQUEST_SIZE_LIMIT.read().await.clone(),
|
|
));
|
|
|
|
let request_size_limit = REQUEST_SIZE_LIMIT.read().await.clone();
|
|
|
|
let cors = CorsLayer::new()
|
|
.allow_methods([http::Method::GET, http::Method::POST, http::Method::DELETE])
|
|
.allow_headers([http::header::CONTENT_TYPE, http::header::AUTHORIZATION])
|
|
.allow_origin(Any);
|
|
|
|
// MCP carries protocol state in its own headers: `MCP-Protocol-Version` from
|
|
// revision 2025-06-18 onward, plus `Mcp-Method` and `Mcp-Name` at 2026-07-28.
|
|
// None of them are CORS-simple, so a browser-based MCP client fails preflight
|
|
// unless they are allowed — hence a separate layer rather than widening the
|
|
// one every other route shares. (`Mcp-Param-*` is only sent for tool inputs
|
|
// annotated with `x-mcp-header`, which no tool here declares.)
|
|
//
|
|
// The request's own header list is mirrored rather than enumerated: a browser
|
|
// MCP client may send any custom name for a preprocessor to read, and no fixed
|
|
// list could cover them. Nothing is granted by echoing it: the origin is
|
|
// `Any`, so browsers never attach credentials, and the endpoint authenticates
|
|
// each request on its own.
|
|
let mcp_cors = CorsLayer::new()
|
|
.allow_methods([http::Method::GET, http::Method::POST, http::Method::DELETE])
|
|
.allow_headers(tower_http::cors::AllowHeaders::mirror_request())
|
|
// The 401 challenge is how a client discovers where to authorize (RFC 9728),
|
|
// and it is not a safelisted response header, so without this a browser
|
|
// client sees an empty one and has no way to begin the OAuth flow.
|
|
.expose_headers([http::header::WWW_AUTHENTICATE])
|
|
.allow_origin(Any);
|
|
|
|
let sp_extension = Arc::new(saml_oss::build_sp_extension().await?);
|
|
|
|
if server_mode {
|
|
#[cfg(feature = "embedding")]
|
|
load_embeddings_db(&db);
|
|
|
|
#[cfg(feature = "cloud")]
|
|
if *CLOUD_HOSTED {
|
|
windmill_queue::init_usage_buffer(db.clone());
|
|
}
|
|
|
|
let mut start_smtp_server = false;
|
|
if let Some(smtp_settings) =
|
|
load_value_from_global_settings(&db, EMAIL_DOMAIN_SETTING).await?
|
|
{
|
|
if smtp_settings.as_str().unwrap_or("") != "" {
|
|
start_smtp_server = true;
|
|
}
|
|
}
|
|
if !start_smtp_server {
|
|
tracing::info!("SMTP server not started because email domain is not set");
|
|
} else {
|
|
#[cfg(feature = "smtp")]
|
|
{
|
|
let smtp_server = Arc::new(SmtpServer {
|
|
db: db.clone(),
|
|
user_db: user_db.clone(),
|
|
auth_cache: auth_cache.clone(),
|
|
base_internal_url: _base_internal_url.clone(),
|
|
});
|
|
let addr = listener
|
|
.local_addr()
|
|
.unwrap_or_else(|_| std::net::SocketAddr::from(([127, 0, 0, 1], 0)));
|
|
if let Err(err) = smtp_server.start_listener_thread(addr).await {
|
|
tracing::error!("Error starting SMTP server: {err:#}");
|
|
}
|
|
}
|
|
#[cfg(not(feature = "smtp"))]
|
|
{
|
|
tracing::info!("SMTP server not started because SMTP feature is not enabled");
|
|
}
|
|
}
|
|
}
|
|
|
|
let job_helpers_service = {
|
|
#[cfg(feature = "parquet")]
|
|
{
|
|
job_helpers_oss::workspaced_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "parquet"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
};
|
|
|
|
// Initialize HTTP trigger refresh loop
|
|
#[cfg(feature = "http_trigger")]
|
|
{
|
|
let http_killpill_rx = killpill_rx.resubscribe();
|
|
// A worker only serves `/r` when one of its own jobs calls the local API, so it loads the
|
|
// routers on that first request rather than at startup.
|
|
triggers::http::refresh_routers_loop(&db, http_killpill_rx, server_mode).await;
|
|
}
|
|
|
|
let triggers_service = triggers::generate_trigger_routers().layer(from_fn_with_state(
|
|
ManageKind::Triggers,
|
|
gate_operator_writes,
|
|
));
|
|
|
|
if !*CLOUD_HOSTED && server_mode && !mcp_mode {
|
|
start_all_listeners(db.clone(), &killpill_rx);
|
|
}
|
|
|
|
if server_mode {
|
|
health::start_health_check_loop(db.clone(), killpill_rx.resubscribe());
|
|
}
|
|
|
|
let port = listener.local_addr().map(|x| x.port()).unwrap_or(8000);
|
|
let ip = listener
|
|
.local_addr()
|
|
.map(|x| x.ip().to_string())
|
|
.unwrap_or("localhost".to_string());
|
|
|
|
// Setup MCP server
|
|
#[allow(unused_variables)]
|
|
let (mcp_router, gateway_mcp_router, mcp_cancellation_token) = {
|
|
#[cfg(feature = "mcp")]
|
|
if server_mode || mcp_mode {
|
|
use mcp::{
|
|
add_www_authenticate_header, add_www_authenticate_header_gateway,
|
|
extract_workspace_from_token, reject_token_query_param,
|
|
};
|
|
let (mcp_router, mcp_cancellation_token) = setup_mcp_server(
|
|
db.clone(),
|
|
user_db,
|
|
_base_internal_url.clone(),
|
|
auth_cache.clone(),
|
|
)
|
|
.await?;
|
|
// Workspace-scoped MCP router
|
|
let workspaced_mcp_router = mcp_router
|
|
.clone()
|
|
.route_layer(from_extractor::<ApiAuthed>())
|
|
.layer(axum::middleware::from_fn(reject_token_query_param))
|
|
.layer(axum::middleware::from_fn(add_www_authenticate_header))
|
|
.layer(axum::middleware::from_fn(extract_and_store_workspace_id));
|
|
// Gateway MCP router — resolves workspace from token
|
|
let gateway_mcp_router = mcp_router
|
|
.route_layer(from_extractor::<ApiAuthed>())
|
|
.layer(axum::middleware::from_fn(extract_workspace_from_token))
|
|
.layer(axum::middleware::from_fn(reject_token_query_param))
|
|
.layer(axum::middleware::from_fn(
|
|
add_www_authenticate_header_gateway,
|
|
));
|
|
(
|
|
workspaced_mcp_router,
|
|
gateway_mcp_router,
|
|
Some(mcp_cancellation_token),
|
|
)
|
|
} else {
|
|
(Router::new(), Router::new(), None)
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
(Router::new(), Router::new(), Option::<()>::None)
|
|
};
|
|
|
|
// Workers block on this before pulling their first job, so it is released ahead of
|
|
// the router tree below. `try_join!` polls this future and `workers_f` on one task,
|
|
// so the yield is what lets them proceed; without it they wait out the whole
|
|
// synchronous build. A request arriving first queues in the bound listener's backlog.
|
|
if let Err(e) = port_tx.send(format!("http://localhost:{}", port)) {
|
|
tracing::error!("Failed to send port: {e:#}");
|
|
return Err(anyhow::anyhow!("Failed to send port, exiting early: {e:#}"));
|
|
}
|
|
tokio::task::yield_now().await;
|
|
|
|
let mcp_list_tools_service = {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
mcp::list_tools_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
};
|
|
|
|
#[cfg(feature = "agent_worker_server")]
|
|
let (agent_workers_router, agent_workers_bg_processor, agent_workers_job_completed_tx) =
|
|
if server_mode {
|
|
windmill_api_agent_workers::workspaced_service(db.clone(), _base_internal_url.clone())
|
|
} else {
|
|
(Router::new(), vec![], None)
|
|
};
|
|
|
|
#[cfg(feature = "agent_worker_server")]
|
|
let agent_cache = Arc::new(AgentCache::new());
|
|
|
|
// build our application with a route
|
|
let app = Router::new()
|
|
.nest(
|
|
"/api",
|
|
Router::new()
|
|
.nest(
|
|
"/w/{workspace_id}",
|
|
Router::new()
|
|
// Reordered alphabetically
|
|
.nest("/acls", granular_acls::workspaced_service())
|
|
// CORS so the opaque-origin in-workspace app viewer (WIN-2006,
|
|
// sandboxed /apps/get) can read the app definition by path
|
|
// (apps/get/p, apps/embed_token/p) with a scoped embed token.
|
|
// Bearer-token-only (no cookies), consistent with the other
|
|
// workspaced services the iframe calls.
|
|
.nest(
|
|
"/apps",
|
|
apps::workspaced_service(request_size_limit * 5).layer(cors.clone()),
|
|
)
|
|
.nest("/assets", windmill_api_assets::workspaced_service())
|
|
.nest("/audit", audit::workspaced_service())
|
|
// Capture configures a trigger without creating one. Its unauthed
|
|
// ingestion routes are a separate service and stay open.
|
|
.nest(
|
|
"/capture",
|
|
capture::workspaced_service().layer(from_fn_with_state(
|
|
ManageKind::Triggers,
|
|
gate_operator_writes,
|
|
)),
|
|
)
|
|
.nest(
|
|
"/concurrency_groups",
|
|
concurrency_groups::workspaced_service(),
|
|
)
|
|
.nest("/dbt", dbt::workspaced_service())
|
|
.nest("/drafts", drafts::workspaced_service())
|
|
.nest("/embeddings", embeddings::workspaced_service())
|
|
.nest("/favorites", favorite::workspaced_service())
|
|
.nest("/flows", flows::workspaced_service())
|
|
.nest("/runnables", runnables::workspaced_service())
|
|
.nest(
|
|
"/workspace_dependencies",
|
|
workspace_dependencies::workspaced_service(),
|
|
)
|
|
// CORS so a chat UI on another origin (an external site, or
|
|
// a sandboxed raw app with its frontend SDK token) can read
|
|
// its conversation history. Bearer-only, like variables.
|
|
.nest(
|
|
"/flow_conversations",
|
|
windmill_api_flow_conversations::workspaced_service()
|
|
.layer(cors.clone()),
|
|
)
|
|
// CORS so an opaque-origin app iframe (WIN-2006 embed,
|
|
// no separate domain) can read folders/listnames with a
|
|
// scoped embed token. Consistent with apps_u/jobs_u cors.
|
|
.nest(
|
|
"/folders",
|
|
folders::workspaced_service().layer(cors.clone()),
|
|
)
|
|
.nest("/folders_history", folder_history::workspaced_service())
|
|
.nest("/groups", groups::workspaced_service())
|
|
.nest("/groups_history", group_history::workspaced_service())
|
|
.nest("/triggers_history", trigger_history::workspaced_service())
|
|
.nest("/inputs", windmill_api_inputs::workspaced_service())
|
|
.nest("/internal_db", internal_db::workspaced_service())
|
|
.route("/labels/list", get(list_workspace_labels))
|
|
.nest("/job_metrics", job_metrics::workspaced_service())
|
|
.nest("/job_helpers", job_helpers_service)
|
|
.nest("/ai_evals", ai_evals::workspaced_service())
|
|
.nest("/jobs", jobs::workspaced_service())
|
|
.nest("/debug", windmill_api_debug::workspaced_service())
|
|
.nest("/native_triggers", {
|
|
#[cfg(feature = "native_trigger")]
|
|
{
|
|
// Only the trigger routes take the gate: the integrations beside
|
|
// them connect the workspace's GitHub/Google account, which is a
|
|
// settings concern and admin-only already.
|
|
native_triggers::handler::generate_native_trigger_routers()
|
|
.layer(from_fn_with_state(
|
|
ManageKind::Triggers,
|
|
gate_operator_writes,
|
|
))
|
|
.merge(
|
|
native_triggers::workspace_integrations::workspaced_service(
|
|
),
|
|
)
|
|
}
|
|
#[cfg(not(feature = "native_trigger"))]
|
|
{
|
|
axum::Router::new()
|
|
}
|
|
})
|
|
.nest("/oauth", {
|
|
#[cfg(feature = "oauth2")]
|
|
{
|
|
oauth2_oss::workspaced_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "oauth2"))]
|
|
Router::new()
|
|
})
|
|
.nest("/mcp/oauth/server", {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
// Only /approve requires authentication (called by frontend)
|
|
mcp::oauth_server::workspaced_authed_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
Router::new()
|
|
})
|
|
.nest("/ai", ai::workspaced_service())
|
|
.nest("/npm_proxy", windmill_api_npm_proxy::workspaced_service())
|
|
.nest(
|
|
"/path_autocomplete",
|
|
path_autocomplete::workspaced_service(),
|
|
)
|
|
.nest("/raw_apps", raw_apps::workspaced_service())
|
|
.nest(
|
|
"/remote_deploy",
|
|
windmill_api_workspaces::remote_deploy::workspaced_service(
|
|
request_size_limit * 5,
|
|
),
|
|
)
|
|
// CORS so the opaque-origin app iframe can read
|
|
// resources/list, resources/type/* with a scoped token.
|
|
.nest(
|
|
"/resources",
|
|
resources::workspaced_service().layer(cors.clone()),
|
|
)
|
|
.nest("/shared_ui", workspace_shared_ui::workspaced_service())
|
|
.nest(
|
|
"/schedules",
|
|
windmill_api_schedule::workspaced_service().layer(from_fn_with_state(
|
|
ManageKind::Schedules,
|
|
gate_operator_writes,
|
|
)),
|
|
)
|
|
.nest("/scripts", scripts::workspaced_service())
|
|
.nest("/trash", trash::workspaced_service())
|
|
.nest(
|
|
"/users",
|
|
// CORS so the opaque-origin app iframe can read
|
|
// users/whoami with a scoped embed token.
|
|
users::workspaced_service()
|
|
.layer(Extension(argon2.clone()))
|
|
.layer(cors.clone()),
|
|
)
|
|
// CORS so a sandboxed raw app's opaque-origin bundle can
|
|
// read variables with its frontend SDK token. Bearer-only:
|
|
// the layer never allows credentials, so no cookie can ride
|
|
// a cross-origin call. Consistent with resources/users.
|
|
.nest(
|
|
"/variables",
|
|
variables::workspaced_service().layer(cors.clone()),
|
|
)
|
|
.nest("/volumes", volumes_oss::workspaced_service())
|
|
.nest("/workers", windmill_api_workers::workspaced_service())
|
|
.nest("/workspaces", workspaces::workspaced_service())
|
|
.nest("/hub", hub_publish::workspaced_service())
|
|
.nest(
|
|
"/data_metrics",
|
|
windmill_api_workspaces::data_metrics::workspaced_service(),
|
|
)
|
|
.nest(
|
|
"/deployment_request",
|
|
windmill_api_workspaces::deployment_requests::workspaced_service(),
|
|
)
|
|
.nest("/oidc", oidc_oss::workspaced_service())
|
|
.nest("/openapi", {
|
|
#[cfg(feature = "http_trigger")]
|
|
{
|
|
windmill_api_openapi::openapi_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "http_trigger"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.merge(triggers_service),
|
|
)
|
|
.nest("/workspaces", workspaces::global_service())
|
|
.nest(
|
|
"/users",
|
|
users::global_service().layer(Extension(argon2.clone())),
|
|
)
|
|
.nest("/settings", windmill_api_settings::global_service())
|
|
.nest("/workers", windmill_api_workers::global_service())
|
|
.nest("/service_logs", service_logs::global_service())
|
|
.nest("/configs", windmill_api_configs::global_service())
|
|
.nest("/scripts", scripts::global_service())
|
|
.nest("/integrations", integration::global_service())
|
|
.nest("/groups", groups::global_service())
|
|
.nest("/flows", flows::global_service())
|
|
.nest("/apps", apps::global_service().layer(cors.clone()))
|
|
.nest("/schedules", windmill_api_schedule::global_service())
|
|
.nest("/embeddings", embeddings::global_service())
|
|
.nest("/ai", ai::global_service())
|
|
.nest("/docs", docs::global_service())
|
|
.nest("/indexer", indexer_oss::management_service())
|
|
.nest("/mcp/w/{workspace_id}/list_tools", mcp_list_tools_service)
|
|
.nest("/db_health", db_health::global_service())
|
|
.nest("/health/detailed", health::detailed_service())
|
|
.nest(
|
|
"/saml",
|
|
saml_oss::authed_service().layer(Extension(Arc::clone(&sp_extension))),
|
|
)
|
|
.nest("/mcp/gateway/oauth/server", {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
mcp::oauth_server::gateway_authed_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
Router::new()
|
|
})
|
|
.route_layer(from_extractor::<ApiAuthed>())
|
|
.route_layer(from_extractor::<users::Tokened>())
|
|
// Workspace-scoped OAuth endpoints that don't require authentication
|
|
// (authorize and token are called by MCP client before user is authenticated)
|
|
.nest("/w/{workspace_id}/mcp/oauth/server", {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
mcp::oauth_server::workspaced_unauthed_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
Router::new()
|
|
})
|
|
// Gateway OAuth endpoints (authorize + token) — no auth required
|
|
.nest("/mcp/gateway/oauth/server", {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
mcp::oauth_server::gateway_unauthed_service().layer(cors.clone())
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
Router::new()
|
|
})
|
|
.nest("/jobs", jobs::global_root_service())
|
|
.nest(
|
|
"/srch/w/{workspace_id}/index",
|
|
indexer_oss::workspaced_service(),
|
|
)
|
|
.nest("/srch/index", indexer_oss::global_service())
|
|
.nest("/oidc", oidc_oss::global_service())
|
|
.nest("/debug", windmill_api_debug::global_service())
|
|
.nest(
|
|
"/saml",
|
|
saml_oss::global_service().layer(Extension(Arc::clone(&sp_extension))),
|
|
)
|
|
.nest(
|
|
"/scim",
|
|
scim_oss::global_service()
|
|
.route_layer(axum::middleware::from_fn(has_scim_token)),
|
|
)
|
|
.nest("/tokens", token::global_service())
|
|
.nest("/concurrency_groups", concurrency_groups::global_service())
|
|
.nest("/scripts_u", scripts::global_unauthed_service())
|
|
.nest("/settings_u", windmill_api_settings::unauthed_service())
|
|
.nest("/apps_u", {
|
|
#[cfg(feature = "enterprise")]
|
|
{
|
|
// CORS so the opaque-origin app viewer (WIN-2006 embed, no
|
|
// separate domain) can load a custom-path public app via
|
|
// public_app_by_custom_path cross-origin. Consistent with
|
|
// the workspaced /w/{workspace_id}/apps_u mount below.
|
|
apps_oss::global_unauthed_service().layer(cors.clone())
|
|
}
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.nest(
|
|
"/w/{workspace_id}/apps_u",
|
|
apps::unauthed_service()
|
|
.layer(from_extractor::<OptAuthed>())
|
|
.layer(cors.clone()),
|
|
)
|
|
.layer(from_extractor::<OptAuthed>())
|
|
// Deprecated, here for backwards compatibility: user should use /mcp/w/{workspace_id}/mcp instead
|
|
.nest(
|
|
"/mcp/w/{workspace_id}/sse",
|
|
mcp_router.clone().layer(mcp_cors.clone()),
|
|
)
|
|
.nest(
|
|
"/mcp/w/{workspace_id}/mcp",
|
|
mcp_router.clone().layer(mcp_cors.clone()),
|
|
)
|
|
.nest("/mcp/gateway", gateway_mcp_router.layer(mcp_cors.clone()))
|
|
.nest("/agent_workers", {
|
|
#[cfg(feature = "agent_worker_server")]
|
|
{
|
|
if let Some(agent_workers_job_completed_tx) =
|
|
agent_workers_job_completed_tx.clone()
|
|
{
|
|
windmill_api_agent_workers::global_service(
|
|
agent_workers_job_completed_tx,
|
|
)
|
|
.layer(Extension(agent_cache.clone()))
|
|
} else {
|
|
Router::new()
|
|
}
|
|
}
|
|
#[cfg(not(feature = "agent_worker_server"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.nest("/w/{workspace_id}/agent_workers", {
|
|
#[cfg(feature = "agent_worker_server")]
|
|
{
|
|
agent_workers_router
|
|
.nest(
|
|
"/volumes",
|
|
volumes_oss::agent_workspaced_service()
|
|
.layer(axum::middleware::from_fn(inject_agent_authed)),
|
|
)
|
|
.layer(Extension(agent_cache.clone()))
|
|
}
|
|
#[cfg(not(feature = "agent_worker_server"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.nest(
|
|
"/w/{workspace_id}/jobs_u",
|
|
jobs::workspace_unauthed_service().layer(cors.clone()),
|
|
)
|
|
.route("/slack", post(slack_approvals::slack_app_callback_handler))
|
|
.nest("/teams", {
|
|
#[cfg(feature = "enterprise")]
|
|
{
|
|
teams_oss::teams_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.route(
|
|
"/w/{workspace_id}/jobs/slack_approval/{job_id}",
|
|
get(slack_approvals::request_slack_approval),
|
|
)
|
|
.route(
|
|
"/w/{workspace_id}/jobs/teams_approval/{job_id}",
|
|
get(teams_approvals_oss::request_teams_approval),
|
|
)
|
|
.nest("/w/{workspace_id}/github_app", {
|
|
#[cfg(feature = "enterprise")]
|
|
{
|
|
git_sync_oss::workspaced_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
Router::new()
|
|
})
|
|
.nest("/github_app", {
|
|
#[cfg(feature = "enterprise")]
|
|
{
|
|
git_sync_oss::global_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
Router::new()
|
|
})
|
|
.nest("/w/{workspace_id}/git_sync", {
|
|
#[cfg(feature = "enterprise")]
|
|
{
|
|
git_sync_oss::workspaced_git_sync_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
Router::new()
|
|
})
|
|
.nest(
|
|
"/w/{workspace_id}/resources_u",
|
|
public_service().layer(cors.clone()),
|
|
)
|
|
.nest(
|
|
"/w/{workspace_id}/capture_u",
|
|
capture::workspaced_unauthed_service().layer(cors.clone()),
|
|
)
|
|
.nest("/w/{workspace_id}/s3_proxy", {
|
|
s3_proxy_oss::workspaced_unauthed_service()
|
|
})
|
|
.nest(
|
|
"/auth",
|
|
users::make_unauthed_service().layer(Extension(argon2)),
|
|
)
|
|
.nest("/oauth", {
|
|
#[cfg(feature = "oauth2")]
|
|
{
|
|
oauth2_oss::global_service().layer(Extension(Arc::clone(&sp_extension)))
|
|
}
|
|
|
|
#[cfg(not(feature = "oauth2"))]
|
|
Router::new()
|
|
})
|
|
.nest("/mcp/oauth", {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
mcp_oauth_oss::global_service()
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
Router::new()
|
|
})
|
|
.nest("/mcp/oauth/server", {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
mcp::oauth_server::global_service().layer(cors.clone())
|
|
}
|
|
|
|
#[cfg(not(feature = "mcp"))]
|
|
Router::new()
|
|
})
|
|
.nest("/r", {
|
|
#[cfg(feature = "http_trigger")]
|
|
{
|
|
triggers::http::handler::http_route_trigger_handler()
|
|
}
|
|
|
|
#[cfg(not(feature = "http_trigger"))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.nest("/gcp/w/{workspace_id}", {
|
|
#[cfg(all(
|
|
feature = "enterprise",
|
|
feature = "gcp_trigger",
|
|
feature = "private"
|
|
))]
|
|
{
|
|
triggers::gcp::handler_oss::gcp_push_route_handler()
|
|
}
|
|
#[cfg(not(all(
|
|
feature = "enterprise",
|
|
feature = "gcp_trigger",
|
|
feature = "private"
|
|
)))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.nest("/azure/w/{workspace_id}", {
|
|
#[cfg(all(
|
|
feature = "enterprise",
|
|
feature = "azure_trigger",
|
|
feature = "private"
|
|
))]
|
|
{
|
|
triggers::azure::handler_oss::azure_push_route_handler()
|
|
}
|
|
#[cfg(not(all(
|
|
feature = "enterprise",
|
|
feature = "azure_trigger",
|
|
feature = "private"
|
|
)))]
|
|
{
|
|
Router::new()
|
|
}
|
|
})
|
|
.route("/version", get(git_v))
|
|
.nest("/health/status", health::status_service())
|
|
.route("/min_keep_alive_version", get(min_keep_alive_version))
|
|
.route("/uptodate", get(is_up_to_date))
|
|
.route("/ee_license", get(ee_license))
|
|
.route("/openapi.yaml", get(openapi))
|
|
.route("/openapi.json", get(openapi_json)),
|
|
)
|
|
// Clients must use workspace-scoped OAuth metadata at:
|
|
// /.well-known/oauth-authorization-server/api/w/{workspace_id}/mcp/oauth/server
|
|
// This is discovered via /.well-known/oauth-protected-resource?workspace_id=...
|
|
.route(
|
|
"/.well-known/oauth-authorization-server/api/w/{workspace_id}/mcp/oauth/server",
|
|
{
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
get(mcp::oauth_server::workspaced_oauth_metadata)
|
|
}
|
|
#[cfg(not(feature = "mcp"))]
|
|
{
|
|
get(|| async { axum::http::StatusCode::NOT_FOUND })
|
|
}
|
|
},
|
|
)
|
|
// RFC 9728 path-based discovery: /.well-known/oauth-protected-resource/api/mcp/w/{workspace_id}/mcp
|
|
.route(
|
|
"/.well-known/oauth-protected-resource/api/mcp/w/{workspace_id}/mcp",
|
|
{
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
get(mcp::oauth_server::protected_resource_metadata_by_path)
|
|
}
|
|
#[cfg(not(feature = "mcp"))]
|
|
{
|
|
get(|| async { axum::http::StatusCode::NOT_FOUND })
|
|
}
|
|
},
|
|
)
|
|
// Gateway OAuth well-known endpoints
|
|
.route(
|
|
"/.well-known/oauth-authorization-server/api/mcp/gateway/oauth/server",
|
|
{
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
get(mcp::oauth_server::gateway_oauth_metadata)
|
|
}
|
|
#[cfg(not(feature = "mcp"))]
|
|
{
|
|
get(|| async { axum::http::StatusCode::NOT_FOUND })
|
|
}
|
|
},
|
|
)
|
|
.route("/.well-known/oauth-protected-resource/api/mcp/gateway", {
|
|
#[cfg(feature = "mcp")]
|
|
{
|
|
get(mcp::oauth_server::gateway_protected_resource_metadata)
|
|
}
|
|
#[cfg(not(feature = "mcp"))]
|
|
{
|
|
get(|| async { axum::http::StatusCode::NOT_FOUND })
|
|
}
|
|
})
|
|
// JWKS endpoint for HashiCorp Vault JWT authentication (must be outside /api prefix)
|
|
.route("/.well-known/jwks.json", {
|
|
#[cfg(all(feature = "private", feature = "enterprise", feature = "openidconnect"))]
|
|
{
|
|
get(crate::oidc_oss::jwks)
|
|
}
|
|
#[cfg(not(all(
|
|
feature = "private",
|
|
feature = "enterprise",
|
|
feature = "openidconnect"
|
|
)))]
|
|
{
|
|
get(windmill_api_settings::get_jwks)
|
|
}
|
|
})
|
|
.fallback(static_assets::static_handler)
|
|
.layer(middleware_stack);
|
|
|
|
let app = if disable_response_logs {
|
|
app
|
|
} else {
|
|
tokio::spawn(async {
|
|
let mut interval =
|
|
tokio::time::interval(std::time::Duration::from_millis(FLUSH_INTERVAL_MS));
|
|
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
|
|
loop {
|
|
interval.tick().await;
|
|
crate::s3_log_batching::flush_s3_batches();
|
|
}
|
|
});
|
|
|
|
app.layer(axum::middleware::from_fn(s3_proxy_log_middleware))
|
|
.layer(
|
|
TraceLayer::new_for_http()
|
|
.on_response(MyOnResponse {})
|
|
.make_span_with(MyMakeSpan {})
|
|
.on_request(())
|
|
.on_failure(MyOnFailure {}),
|
|
)
|
|
};
|
|
|
|
let app = if let Some(domain) = public_app_layer::PUBLIC_APP_DOMAIN.as_ref() {
|
|
tracing::info!("Public app domain filter enabled for domain: {}", domain);
|
|
app.layer(axum::middleware::from_fn(
|
|
public_app_layer::public_app_domain_filter,
|
|
))
|
|
} else {
|
|
app
|
|
};
|
|
|
|
// Seed the per-request LogContext task-local. Registered outside
|
|
// TraceLayer so MyOnResponse::on_response's `"response"` log fires inside
|
|
// the scope and gets method/uri/workspace_id/email attached by the EE
|
|
// LogContextBridge.
|
|
let app = app.layer(axum::middleware::from_fn(
|
|
tracing_init::log_context_middleware,
|
|
));
|
|
|
|
let app = app.layer(axum::middleware::from_fn(set_deploy_origin));
|
|
|
|
let app = app.layer(axum::middleware::from_fn(set_request_client));
|
|
|
|
let app = app.layer(CatchPanicLayer::custom(|err| {
|
|
tracing::error!("panic in handler, returning 500: {:?}", err);
|
|
Response::builder()
|
|
.status(http::StatusCode::INTERNAL_SERVER_ERROR)
|
|
.body(Body::from("Internal Server Error"))
|
|
.unwrap()
|
|
}));
|
|
|
|
if let Some(name) = name.as_ref() {
|
|
tracing::info!("server starting for name={name}");
|
|
}
|
|
let listener = listener.tap_io(move |tcp_stream| {
|
|
let _ = tcp_stream.set_nodelay(!server_mode);
|
|
});
|
|
let server = axum::serve(listener, app.into_make_service());
|
|
|
|
tracing::info!(
|
|
instance = %*INSTANCE_NAME,
|
|
"server started on port={} and addr={} {}",
|
|
port,
|
|
ip,
|
|
name.map(|x| format!("name={x}")).unwrap_or_default()
|
|
);
|
|
|
|
// Announce this server is ready so coordinated restarts can detect a healthy peer.
|
|
if server_mode {
|
|
if let Err(e) = announce_server_started(&db).await {
|
|
tracing::warn!("Failed to announce server started: {e:#}");
|
|
}
|
|
}
|
|
|
|
let server = server.with_graceful_shutdown(async move {
|
|
killpill_rx.recv().await.ok();
|
|
#[cfg(feature = "agent_worker_server")]
|
|
if let Some(agent_workers_job_completed_tx) = agent_workers_job_completed_tx {
|
|
if let Err(e) = agent_workers_job_completed_tx.kill().await {
|
|
tracing::error!("Error killing agent workers: {e:#}");
|
|
}
|
|
}
|
|
tracing::info!("Graceful shutdown of server");
|
|
|
|
#[cfg(feature = "mcp")]
|
|
if let Some(mcp_cancellation_token) = mcp_cancellation_token {
|
|
mcp_cancellation_token.cancel();
|
|
tracing::info!("MCP server shutdown");
|
|
}
|
|
});
|
|
|
|
server.await?;
|
|
#[cfg(feature = "agent_worker_server")]
|
|
for (i, bg_processor) in agent_workers_bg_processor.into_iter().enumerate() {
|
|
tracing::info!("server off. shutting down agent worker bg processor {i}");
|
|
bg_processor.await?;
|
|
tracing::info!("agent worker bg processor {i} shut down");
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
async fn is_up_to_date() -> Result<String, AppError> {
|
|
let error_reading_version = || anyhow::anyhow!("Error reading latest released version");
|
|
let version = HTTP_CLIENT
|
|
.get("https://api.github.com/repos/windmill-labs/windmill/releases/latest")
|
|
.timeout(Duration::from_secs(10))
|
|
.send()
|
|
.await
|
|
.context("Impossible to reach api.github")?
|
|
.json::<serde_json::Value>()
|
|
.await?
|
|
.get("tag_name")
|
|
.ok_or_else(error_reading_version)?
|
|
.as_str()
|
|
.ok_or_else(error_reading_version)?
|
|
.to_string();
|
|
let release = GIT_VERSION
|
|
.split('-')
|
|
.next()
|
|
.ok_or_else(error_reading_version)?
|
|
.to_string();
|
|
if version == release {
|
|
Ok("yes".to_string())
|
|
} else {
|
|
Ok(format!("Update: {GIT_VERSION} -> {version}"))
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "enterprise")]
|
|
async fn git_v() -> String {
|
|
format!("EE {GIT_VERSION}")
|
|
}
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
async fn git_v() -> String {
|
|
format!("CE {GIT_VERSION}")
|
|
}
|
|
|
|
async fn min_keep_alive_version() -> Json<serde_json::Value> {
|
|
let worker = windmill_common::min_version::MIN_KEEP_ALIVE_VERSION;
|
|
let agent = windmill_common::min_version::AGENT_MIN_KEEP_ALIVE_VERSION;
|
|
Json(serde_json::json!({
|
|
"worker": format!("{}.{}.{}", worker.0, worker.1, worker.2),
|
|
"agent": format!("{}.{}.{}", agent.0, agent.1, agent.2)
|
|
}))
|
|
}
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
async fn ee_license() -> &'static str {
|
|
""
|
|
}
|
|
|
|
async fn list_workspace_labels(
|
|
Extension(db): Extension<DB>,
|
|
axum::extract::Path(w_id): axum::extract::Path<String>,
|
|
) -> windmill_common::error::JsonResult<Vec<String>> {
|
|
let labels = sqlx::query_scalar!(
|
|
"SELECT DISTINCT unnest(labels) as \"label!\" FROM (
|
|
SELECT labels FROM script WHERE workspace_id = $1 AND labels IS NOT NULL
|
|
UNION ALL SELECT labels FROM flow WHERE workspace_id = $1 AND labels IS NOT NULL
|
|
UNION ALL SELECT labels FROM resource WHERE workspace_id = $1 AND labels IS NOT NULL
|
|
UNION ALL SELECT labels FROM variable WHERE workspace_id = $1 AND labels IS NOT NULL
|
|
UNION ALL SELECT labels FROM schedule WHERE workspace_id = $1 AND labels IS NOT NULL
|
|
UNION ALL SELECT labels FROM app WHERE workspace_id = $1 AND labels IS NOT NULL
|
|
UNION ALL SELECT labels FROM folder WHERE workspace_id = $1 AND labels IS NOT NULL
|
|
) t ORDER BY 1",
|
|
&w_id
|
|
)
|
|
.fetch_all(&db)
|
|
.await?;
|
|
Ok(axum::Json(labels))
|
|
}
|
|
|
|
#[cfg(feature = "enterprise")]
|
|
async fn ee_license() -> String {
|
|
use windmill_common::ee_oss::{LICENSE_KEY_ID, LICENSE_KEY_VALID};
|
|
|
|
if LICENSE_KEY_VALID.load(std::sync::atomic::Ordering::Relaxed) {
|
|
(**LICENSE_KEY_ID.load()).clone()
|
|
} else {
|
|
"".to_string()
|
|
}
|
|
}
|
|
|
|
async fn openapi() -> Response {
|
|
Response::builder()
|
|
.header("content-type", "application/yaml")
|
|
.body(Body::from(include_str!("../openapi-deref.yaml")))
|
|
.unwrap()
|
|
}
|
|
|
|
async fn openapi_json() -> Response {
|
|
Response::builder()
|
|
.header("content-type", "application/json")
|
|
.body(Body::from(include_str!("../openapi-deref.json")))
|
|
.unwrap()
|
|
}
|
|
|
|
pub async fn migrate_db(
|
|
db: &DB,
|
|
killpill_rx: tokio::sync::broadcast::Receiver<()>,
|
|
) -> anyhow::Result<Option<JoinHandle<()>>> {
|
|
db::migrate(db, killpill_rx)
|
|
.await
|
|
.map_err(|e| anyhow::anyhow!("Error migrating db: {e:#}"))
|
|
}
|
|
|
|
pub async fn wait_for_db_migrations(
|
|
db: &DB,
|
|
killpill_rx: tokio::sync::broadcast::Receiver<()>,
|
|
) -> anyhow::Result<()> {
|
|
db::wait_for_migrations(db, killpill_rx)
|
|
.await
|
|
.map_err(|e| anyhow::anyhow!("Error waiting for db migrations: {e:#}"))
|
|
}
|
|
|
|
const SERVER_HEARTBEAT_TASK: &str = "server_heartbeat";
|
|
|
|
/// Write a server-started heartbeat to `background_task_state` so that other
|
|
/// traffic-serving instances waiting to restart can detect this one is healthy.
|
|
///
|
|
/// Only `server_mode` processes announce, since only they are peers worth waiting for:
|
|
/// `spawn_graceful_killpill` holds a shutdown open to keep the API answered, and a worker,
|
|
/// indexer or MCP process coming up is no evidence that it is.
|
|
///
|
|
/// The row is keyed per host and `owner` per process, and both halves carry weight.
|
|
/// `INSTANCE_NAME` is random per start, so a row keyed on it never conflicts and
|
|
/// accumulates one row per start; `owner` is what tells a peer's start from its own
|
|
/// when processes share a host.
|
|
async fn announce_server_started(db: &DB) -> anyhow::Result<()> {
|
|
use windmill_common::{utils::HOSTNAME, INSTANCE_NAME};
|
|
|
|
let instance = INSTANCE_NAME.as_str();
|
|
let host = HOSTNAME.as_str();
|
|
|
|
sqlx::query(
|
|
"INSERT INTO background_task_state (name, value, running, owner, started_at, updated_at)
|
|
VALUES ($1, '\"started\"'::jsonb, true, $2, NOW(), NOW())
|
|
ON CONFLICT (name)
|
|
DO UPDATE SET started_at = NOW(), updated_at = NOW(), running = true, owner = $2",
|
|
)
|
|
.bind(format!("{SERVER_HEARTBEAT_TASK}:{host}"))
|
|
.bind(instance)
|
|
.execute(db)
|
|
.await?;
|
|
|
|
tracing::info!("Announced server started for instance {instance} on host {host}");
|
|
Ok(())
|
|
}
|
|
|
|
/// Check whether any server instance (other than ourselves) has announced
|
|
/// itself as started after `not_before` (i.e. after the restart was initiated).
|
|
pub async fn check_any_server_started(db: &DB, not_before: chrono::DateTime<chrono::Utc>) -> bool {
|
|
use windmill_common::INSTANCE_NAME;
|
|
|
|
let my_instance = INSTANCE_NAME.as_str();
|
|
let prefix = format!("{SERVER_HEARTBEAT_TASK}:");
|
|
|
|
sqlx::query_scalar!(
|
|
"SELECT EXISTS(
|
|
SELECT 1 FROM background_task_state
|
|
WHERE name LIKE $1
|
|
AND owner != $2
|
|
AND running = true
|
|
AND updated_at > $3
|
|
) AS \"exists!\"",
|
|
format!("{prefix}%"),
|
|
my_instance,
|
|
not_before
|
|
)
|
|
.fetch_one(db)
|
|
.await
|
|
.unwrap_or(false)
|
|
}
|
|
|
|
/// Delete `server_heartbeat:*` rows that have not been refreshed in a long
|
|
/// time. A restart in place reuses its host's row (see
|
|
/// `announce_server_started`), but a host that never comes back leaves one
|
|
/// behind, and hosts are disposable wherever the hostname carries a generated
|
|
/// pod or container id.
|
|
///
|
|
/// The row is only consulted by `check_any_server_started`, which itself
|
|
/// filters on `updated_at > not_before` (the moment a restart was initiated),
|
|
/// so rows older than the cutoff cannot influence any restart decision and
|
|
/// are safe to delete.
|
|
pub async fn cleanup_stale_server_heartbeats(db: &DB) -> anyhow::Result<u64> {
|
|
let prefix = format!("{SERVER_HEARTBEAT_TASK}:");
|
|
let res = sqlx::query!(
|
|
"DELETE FROM background_task_state
|
|
WHERE name LIKE $1
|
|
AND updated_at < NOW() - INTERVAL '7 days'",
|
|
format!("{prefix}%"),
|
|
)
|
|
.execute(db)
|
|
.await?;
|
|
Ok(res.rows_affected())
|
|
}
|