Files
windmill/backend/windmill-api/src/lib.rs
T
Diego Imbert 386e115daa feat: lazily expand s3 explorer folders one level at a time (#10420)
* feat: wire paged object storage listing module

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* feat: document list_stored_files_paged endpoint in openapi

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* feat: lazily expand s3 explorer folders one level at a time

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: pin ee-repo-ref to the paged listing branch

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: share object_store credential resolution and surface listing errors

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sqq2LhmWaGwP11Cf3UqWxe

* Chevron is cool

* page size 5000

* feat: make the load more row full-width, secondary and chevron-led

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sqq2LhmWaGwP11Cf3UqWxe

* fix: render newly loaded flat pages inside already-expanded folders

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: address review findings in the lazy s3 explorer

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: address review nits in the lazy s3 explorer

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: bump ee-repo-ref after merging main

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: document ambient credential contract and constrain max_keys schema

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: treat an exhausted page token as exhausted, not as a continuation

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: bump ee-repo-ref for canonical prefix validation

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: bump ee-repo-ref for prefix scoping and opaque cursors

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: invalidate a folder's in-flight load when deleting from it

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: discard a stale folder page after its level is invalidated

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: bump ee-repo-ref for bounded local listing

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: label folders whose final path segment is empty

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* feat: search files by any part of their path, not just folder prefix

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* feat: search files by path prefix instead of a full-bucket scan

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* fix: guard stale search responses and describe prefix search accurately

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: bump ee-repo-ref for the search prefix fallback fix

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: regenerate the served openapi specs

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: bump ee-repo-ref for the search cursor fallback fix

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JqAc8mz6Gu698kBbJJVwcT

* chore: bump ee-repo-ref for the bounded search scan

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: surface a failed flat listing instead of spinning forever

The flat branch of loadFiles was awaited without a catch, and loadFlatFiles
clears its loading flags only on the success tail. Every caller reaches it
un-awaited, so a rejected listing left the drawer on "Loading content" with
nothing reported. Routing the filter box through this arm made it reachable
per keystroke rather than once per open.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: give back the flat cursor when a page fails to load

"Load more" advanced `page` before requesting it, so a failed page left the
cursor pointing at a `listMarkers` slot that was never filled. The retry sent
no marker at all and silently replayed the first page, and the
`listMarkers.length == page` guard kept it there until the listing was reset.

Only reachable now that a failed page is retryable rather than a permanent
spinner.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: scope the flat cursor rollback to its own listing

The rollback matched on the page number alone, so a page that failed after a
filter or storage change could roll back the *replacement* listing once it had
reached the same number, stranding its cursor. Tie it to the generation the
request was issued under.

The delete replay loop had the mirrored problem: it re-drove `page` by hand and
carried on past a failed page, leaving `page` ahead of `listMarkers` for good.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: skip the delete replay when the fresh listing itself failed

clearAndLoadFiles dropped the result it already computes, so a failed
post-delete listing still ran the replay loop: each page advanced `page` with
an empty `listMarkers`, which never recovers because the marker-length guard
only pushes when the two agree. Every later "Load more" then replayed page one.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: stop a superseded lazy load from writing into the search that replaced it

loadFolderPage resolves rather than throwing once its generation is stale, so a
filter change that switches the picker to the flat listing mid-flight left the
lazy branch free to expand a preselected file into the search's results and to
clear the search's loading flags. Guard both on the generation it started under,
as the flat branch already does.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: check the listing generation throughout the reveal walk

Revealing a preselected key is a chain of round trips, so checking once at entry
left the rest of the walk free to keep loading after a filter change had already
switched the picker to the search — under the replacement generation, so the
per-level guards inside loadFolderPage saw nothing wrong.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: let a late metadata failure clear only its own preview

The handler blanked fileMetadata and filePreview without checking that its
request still owned the pane, so selecting a second file while the first was
still loading meant the first's rejection wiped the second's preview.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: key preview ownership on the request, not the selected key

Comparing the selected key let an older request speak for a newer one when both
targeted the same key, which a storage switch does, and made a request whose
selection had moved to something with no metadata return early with the spinner
still up — the case the handler exists to prevent.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: clear the preview when the previewed file is deleted

The lazy branch refetches only the affected level and returns, so it never
reached the reset that the flat refresh gets from clearAndLoadFiles. The pane
renders from fileMetadata rather than from the selection, leaving the deleted
file previewed with working download, move and delete actions.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: retire the in-flight preview load when its file is deleted

Clearing the pane was not enough: a metadata response computed before the DELETE
landed still repopulated it, restoring the deleted file's preview and its
download, move and delete actions. Deleting now retires the owning request, and
the success and preview writes honour that the same way the failure path does.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: clear the preview loading flag when the delete retires its request

Retiring the in-flight metadata load left nobody to report its outcome, so in
lazy mode the pane sat on "Loading..." instead of falling back to the empty
state. The delete owns the flag once it has retired the request.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* chore: drop the regenerated openapi deref artifacts

They are generated files that CI only syntax-validates, never checks against
openapi.yaml, and the committed copies already differ from the spec they derive
from by ~9.7k lines. Regenerating here imported that pre-existing drift into a
feature diff, burying ~800 lines of actual change under ~17k lines of other
changes' staleness. Regenerating them is its own chore.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* docs: state the flat cursor invariant once, where the cursor lives

It was spelled out at four sites, which is what AGENTS.md asks not to do. The
rule now sits on the declaration it constrains and the guards reference it.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* nit ui

* fix: add the paged listing to the served openapi json

openapi_json() embeds openapi-deref.json via include_str!, and the Docker build
regenerates only the yaml artifact, so the json is served exactly as committed —
leaving the new operation out of the Scalar API reference.

Spliced in the operation and the two schemas it references rather than
regenerating, which would have re-imported ~7k lines of pre-existing drift
between the committed artifact and the spec it derives from.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* chore: bump ee-repo-ref for the filesystem symlink boundary

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* chore: update ee-repo-ref to 0373b4bfdaf8dd51533552e2e4de63ceb3c18b4d

This commit updates the EE repository reference after PR #697 was merged in windmill-ee-private.

Previous ee-repo-ref: eb1a765bb9b29e0c94a6e4942c304934fa15406e

New ee-repo-ref: 0373b4bfdaf8dd51533552e2e4de63ceb3c18b4d

Automated by sync-ee-ref workflow.

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
2026-08-04 01:41:06 +02:00

1364 lines
50 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;
#[cfg(feature = "enterprise")]
use windmill_api_auth::ee_oss::ExternalJwks;
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;
mod ai_skills;
mod apps;
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 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 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;
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
}
#[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,
});
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);
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();
triggers::http::refresh_routers_loop(&db, http_killpill_rx).await;
}
let triggers_service = triggers::generate_trigger_routers();
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,
};
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(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(
add_www_authenticate_header_gateway,
))
.layer(axum::middleware::from_fn(extract_workspace_from_token));
(
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)
};
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())
.nest("/capture", capture::workspaced_service())
.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(),
)
.nest(
"/flow_conversations",
windmill_api_flow_conversations::workspaced_service(),
)
// 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("/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("/jobs", jobs::workspaced_service())
.nest("/debug", windmill_api_debug::workspaced_service())
.nest("/native_triggers", {
#[cfg(feature = "native_trigger")]
{
native_triggers::handler::generate_native_trigger_routers().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("/ai_skills", ai_skills::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())
// 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())
.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(cors.clone()),
)
.nest(
"/mcp/w/{workspace_id}/mcp",
mcp_router.clone().layer(cors.clone()),
)
.nest("/mcp/gateway", gateway_mcp_router.layer(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}/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(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()
);
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:#}"));
}
// Announce this server is ready so coordinated restarts can detect a healthy peer.
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 instances waiting to restart can detect this server is healthy.
async fn announce_server_started(db: &DB) -> anyhow::Result<()> {
use windmill_common::INSTANCE_NAME;
let instance = INSTANCE_NAME.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 updated_at = NOW(), running = true, owner = $2",
)
.bind(format!("{SERVER_HEARTBEAT_TASK}:{instance}"))
.bind(instance)
.execute(db)
.await?;
tracing::info!("Announced server started for instance {instance}");
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. Each server startup generates a fresh random `INSTANCE_NAME` and
/// inserts a new row keyed by `server_heartbeat:{instance}`; because that
/// row is only written once (on startup) and never updated thereafter, the
/// table grows by one row per server restart and is otherwise never pruned.
///
/// 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())
}