From c2c1be4a486a21d05d511a552180e199933bdef4 Mon Sep 17 00:00:00 2001 From: centdix <40307056+centdix@users.noreply.github.com> Date: Wed, 23 Apr 2025 23:28:41 +0200 Subject: [PATCH] feat: Add MCP endpoints (#5639) --- backend/Cargo.lock | 89 ++- backend/Cargo.toml | 1 + backend/windmill-api/Cargo.toml | 2 + backend/windmill-api/src/jobs.rs | 1 - backend/windmill-api/src/lib.rs | 52 +- backend/windmill-api/src/mcp.rs | 642 ++++++++++++++++++ .../src/lib/components/TokensTable.svelte | 295 ++++++++ .../src/lib/components/UserSettings.svelte | 222 ++---- .../details/EmailTriggerConfigSection.svelte | 7 +- .../webhook/WebhooksConfigSection.svelte | 17 +- .../src/routes/(root)/(logged)/+layout.svelte | 7 +- .../svix/create-webhook/+page@(root).svelte | 2 +- .../user/(user)/workspaces/+page.svelte | 22 +- 13 files changed, 1102 insertions(+), 257 deletions(-) create mode 100644 backend/windmill-api/src/mcp.rs create mode 100644 frontend/src/lib/components/TokensTable.svelte diff --git a/backend/Cargo.lock b/backend/Cargo.lock index c0256d1816..4f3ca231e6 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -6081,7 +6081,7 @@ dependencies = [ "js-sys", "log", "wasm-bindgen", - "windows-core 0.61.0", + "windows-core 0.57.0", ] [[package]] @@ -6798,7 +6798,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fc2f4eb4bc735547cfed7c0a4922cbd04a4655978c09b54f1f7b228750664c34" dependencies = [ "cfg-if", - "windows-targets 0.52.6", + "windows-targets 0.48.5", ] [[package]] @@ -9870,6 +9870,40 @@ dependencies = [ "syn 1.0.109", ] +[[package]] +name = "rmcp" +version = "0.1.5" +source = "git+https://github.com/windmill-labs/rust-sdk#e6c368965711a2afe45218149d1be5324b098a4d" +dependencies = [ + "async-stream", + "axum", + "base64 0.21.7", + "chrono", + "futures", + "paste", + "pin-project-lite", + "rand 0.9.0", + "rmcp-macros", + "schemars", + "serde", + "serde_json", + "thiserror 2.0.12", + "tokio", + "tokio-stream", + "tokio-util", + "tracing", +] + +[[package]] +name = "rmcp-macros" +version = "0.1.5" +source = "git+https://github.com/windmill-labs/rust-sdk#e6c368965711a2afe45218149d1be5324b098a4d" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.100", +] + [[package]] name = "ron" version = "0.8.1" @@ -13771,7 +13805,7 @@ dependencies = [ "log", "naga", "once_cell", - "parking_lot 0.12.3", + "parking_lot 0.11.2", "profiling", "raw-window-handle", "ron", @@ -13813,7 +13847,7 @@ dependencies = [ "ndk-sys", "objc", "once_cell", - "parking_lot 0.12.3", + "parking_lot 0.11.2", "profiling", "range-alloc", "raw-window-handle", @@ -14020,6 +14054,7 @@ dependencies = [ "rdkafka", "regex", "reqwest 0.12.15", + "rmcp", "rsa", "rumqttc", "rust-embed", @@ -14045,7 +14080,7 @@ dependencies = [ "tokio-tungstenite", "tokio-util", "tonic", - "tower 0.5.2", + "tower 0.4.13", "tower-cookies", "tower-http", "tracing", @@ -14624,19 +14659,6 @@ dependencies = [ "windows-targets 0.52.6", ] -[[package]] -name = "windows-core" -version = "0.61.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4763c1de310c86d75a878046489e2e5ba02c649d185f21c67d4cf8a56d098980" -dependencies = [ - "windows-implement 0.60.0", - "windows-interface 0.59.1", - "windows-link", - "windows-result 0.3.2", - "windows-strings 0.4.0", -] - [[package]] name = "windows-implement" version = "0.56.0" @@ -14670,17 +14692,6 @@ dependencies = [ "syn 2.0.100", ] -[[package]] -name = "windows-implement" -version = "0.60.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a47fddd13af08290e67f4acabf4b459f647552718f683a7b415d290ac744a836" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.100", -] - [[package]] name = "windows-interface" version = "0.56.0" @@ -14714,17 +14725,6 @@ dependencies = [ "syn 2.0.100", ] -[[package]] -name = "windows-interface" -version = "0.59.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bd9211b69f8dcdfa817bfd14bf1c97c9188afa36f4750130fcdf3f400eca9fa8" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.100", -] - [[package]] name = "windows-link" version = "0.1.1" @@ -14788,15 +14788,6 @@ dependencies = [ "windows-link", ] -[[package]] -name = "windows-strings" -version = "0.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7a2ba9642430ee452d5a7aa78d72907ebe8cfda358e8cb7918a2050581322f97" -dependencies = [ - "windows-link", -] - [[package]] name = "windows-sys" version = "0.48.0" diff --git a/backend/Cargo.toml b/backend/Cargo.toml index e1f4403506..ec8ab26f84 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -72,6 +72,7 @@ dind = ["windmill-worker/dind"] websocket = ["windmill-api/websocket"] http_trigger = ["windmill-api/http_trigger"] postgres_trigger = ["windmill-api/postgres_trigger"] +mcp = ["windmill-api/mcp"] mqtt_trigger = ["windmill-api/mqtt_trigger"] sqs_trigger = ["windmill-api/sqs_trigger", "windmill-common/aws_auth", "windmill-api/openidconnect"] gcp_trigger = ["windmill-api/gcp_trigger"] diff --git a/backend/windmill-api/Cargo.toml b/backend/windmill-api/Cargo.toml index c6702b8ac3..404eb4660e 100644 --- a/backend/windmill-api/Cargo.toml +++ b/backend/windmill-api/Cargo.toml @@ -35,8 +35,10 @@ sqs_trigger = ["dep:aws-sdk-sqs", "dep:thiserror", "dep:aws-config"] deno_core = ["dep:deno_core", "dep:deno_error"] gcp_trigger = ["dep:thiserror", "dep:google-cloud-pubsub", "dep:google-cloud-googleapis", "dep:tonic"] cloud = ["windmill-common/cloud"] +mcp = ["dep:rmcp"] [dependencies] +rmcp = { git = "https://github.com/windmill-labs/rust-sdk", features = ["transport-sse-server"], optional = true } windmill-queue.workspace = true windmill-common = { workspace = true, default-features = false } windmill-audit.workspace = true diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 1ed0871a7c..7913796441 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -4622,7 +4622,6 @@ pub async fn run_wait_result_script_by_hash( check_license_key_valid().await?; let args = args.to_push_args_owned(&authed, &db, &w_id).await?; - check_queue_too_long(&db, run_query.queue_limit).await?; let hash = script_hash.0; diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index ece80c64b1..d773cc8615 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -18,6 +18,8 @@ use crate::oauth2_ee::SlackVerifier; #[cfg(feature = "smtp")] use crate::smtp_server_ee::SmtpServer; +#[cfg(feature = "mcp")] +use crate::mcp::{setup_mcp_server, Runner as McpRunner}; use crate::tracing_init::MyOnFailure; use crate::{ tracing_init::{MyMakeSpan, MyOnResponse}, @@ -27,6 +29,7 @@ use crate::{ #[cfg(feature = "agent_worker_server")] use agent_workers_ee::AgentCache; + use anyhow::Context; use argon2::Argon2; use axum::extract::DefaultBodyLimit; @@ -139,6 +142,9 @@ mod workspaces_ee; mod workspaces_export; mod workspaces_extra; +#[cfg(feature = "mcp")] +mod mcp; + pub const DEFAULT_BODY_LIMIT: usize = 2097152 * 100; // 200MB lazy_static::lazy_static! { @@ -448,6 +454,23 @@ pub async fn run_server( } } + let listener = tokio::net::TcpListener::bind(addr) + .await + .context("binding main windmill server")?; + 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 + #[cfg(feature = "mcp")] + let (mcp_sse_server, mcp_router) = setup_mcp_server(addr, "/api/w/:workspace_id/mcp")?; + #[cfg(feature = "mcp")] + let mcp_main_ct = mcp_sse_server.config.ct.clone(); // Token to signal shutdown *to* MCP + #[cfg(feature = "mcp")] + let mcp_service_ct = mcp_sse_server.with_service(McpRunner::new); // Token to wait for MCP *service* shutdown + #[cfg(feature = "agent_worker_server")] let (agent_workers_router, agent_workers_bg_processor, agent_workers_killpill_tx) = agent_workers_ee::workspaced_service(db.clone(), _base_internal_url.clone()); @@ -586,6 +609,17 @@ pub async fn run_server( .layer(from_extractor::()) .layer(cors.clone()), ) + .nest("/w/:workspace_id/mcp", { + #[cfg(feature = "mcp")] + { + mcp_router + } + #[cfg(not(feature = "mcp"))] + { + Router::new() + } + }) + .layer(from_extractor::()) .nest( "/w/:workspace_id/jobs_u", jobs::workspace_unauthed_service().layer(cors.clone()), @@ -694,14 +728,6 @@ pub async fn run_server( .on_failure(MyOnFailure {}), ) }; - let listener = tokio::net::TcpListener::bind(addr) - .await - .context("binding main windmill server")?; - 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()); let server = axum::serve(listener, app.into_make_service()); @@ -723,10 +749,18 @@ pub async fn run_server( tracing::error!("Error killing agent workers: {e:#}"); } tracing::info!("Graceful shutdown of server"); + + #[cfg(feature = "mcp")] + { + tracing::info!("Received shutdown signal, cancelling MCP server..."); + mcp_main_ct.cancel(); + tracing::info!("Waiting for MCP service cancellation..."); + mcp_service_ct.cancelled().await; + tracing::info!("MCP service cancelled."); + } }); 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}"); diff --git a/backend/windmill-api/src/mcp.rs b/backend/windmill-api/src/mcp.rs new file mode 100644 index 0000000000..b30f80d643 --- /dev/null +++ b/backend/windmill-api/src/mcp.rs @@ -0,0 +1,642 @@ +use std::borrow::Cow; +use std::collections::HashMap; +use std::net::SocketAddr; +use std::sync::Arc; + +use axum::body::to_bytes; +use axum::Router; +use rmcp::transport::sse_server::{SseServer, SseServerConfig}; +use rmcp::{ + handler::server::ServerHandler, + model::*, + service::{RequestContext, RoleServer}, + Error, +}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use sql_builder::prelude::*; +use sqlx::FromRow; +use tokio::try_join; +use tokio_util::sync::CancellationToken; +use windmill_common::db::UserDB; +use windmill_common::worker::to_raw_value; +use windmill_common::DB; + +use crate::db::ApiAuthed; +use crate::jobs::{ + run_wait_result_flow_by_path_internal, run_wait_result_script_by_path_internal, RunJobQuery, +}; +use windmill_common::utils::StripPath; + +#[derive(Clone)] +pub struct Runner {} + +#[derive(Serialize, Deserialize, Debug, Clone, Default)] +struct Schema { + #[serde(default)] + properties: HashMap, + #[serde(flatten)] + other: HashMap, +} + +#[derive(Serialize, Deserialize, Debug, Clone, Default)] +struct SchemaProperty { + #[serde(skip_serializing_if = "Option::is_none")] + format: Option, + #[serde(skip_serializing_if = "Option::is_none")] + description: Option, + #[serde(skip_serializing_if = "Option::is_none")] + r#type: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[allow(non_snake_case)] + oneOf: Option>, + #[serde(flatten)] + other: HashMap, +} + +#[derive(Serialize, FromRow)] +struct ScriptInfo { + path: String, + summary: Option, + description: Option, + schema: Option, +} + +#[derive(Serialize, FromRow, Debug)] +struct FlowInfo { + path: String, + summary: Option, + description: Option, + schema: Option, +} + +#[derive(Serialize, FromRow, Debug)] +struct ResourceInfo { + path: String, + description: Option, + resource_type: String, +} + +#[derive(Serialize, FromRow, Debug)] +struct ResourceType { + name: String, + description: Option, +} + +#[derive(Serialize, FromRow, Debug)] +struct ResourceCache { + resource_type: ResourceType, + resources: Vec, +} + +impl Runner { + pub fn new() -> Self { + Self {} + } + + fn transform_path(path: &str, type_str: &str) -> Result { + if type_str != "script" && type_str != "flow" { + return Err(format!("Invalid type: {}", type_str)); + } + + // Only apply special underscore escaping for paths starting with "f/" + let transformed = if path.starts_with("f/") { + let escaped_path = path.replace('_', "__"); + escaped_path.replace('/', "_") + } else { + path.replace('/', "_") + }; + + Ok(format!("{}-{}", type_str, transformed)) + } + + fn reverse_transform(transformed_path: &str) -> Result<(&str, String), String> { + let prefix = if transformed_path.starts_with("script-") { + "script-" + } else if transformed_path.starts_with("flow-") { + "flow-" + } else { + return Err(format!( + "Invalid prefix in transformed path: {}", + transformed_path + )); + }; + + let type_str = &prefix[..prefix.len() - 1]; // "script" or "flow" + let mangled_path = &transformed_path[prefix.len()..]; + + // Check if this path was previously transformed with special underscore handling + let is_special_path = mangled_path.starts_with("f_"); + + let original_path = if is_special_path { + const TEMP_PLACEHOLDER: &str = "@@UNDERSCORE@@"; + let path_with_placeholder = mangled_path.replace("__", TEMP_PLACEHOLDER); + let path_with_slashes = path_with_placeholder.replace('_', "/"); + path_with_slashes.replace(TEMP_PLACEHOLDER, "_") + } else { + mangled_path.replacen('_', "/", 2) + }; + + Ok((type_str, original_path)) + } + + async fn inner_get_resource_type_info( + user_db: &UserDB, + authed: &ApiAuthed, + workspace_id: &str, + resource_type: &str, + ) -> Result { + let mut sqlb = SqlBuilder::select_from("resource_type as o"); + sqlb.fields(&["o.name", "o.description"]); + sqlb.and_where("o.workspace_id = ?".bind(&workspace_id)); + sqlb.and_where("o.name = ?".bind(&resource_type)); + let sql = sqlb.sql().map_err(|_e| { + tracing::error!("failed to build sql: {}", _e); + Error::internal_error("failed to build sql", None) + })?; + let mut tx = user_db + .clone() + .begin(authed) + .await + .map_err(|_e| Error::internal_error("failed to begin transaction", None))?; + let rows = sqlx::query_as::<_, ResourceType>(&sql) + .fetch_one(&mut *tx) + .await + .map_err(|_e| { + tracing::error!("Failed to fetch resource info: {}", _e); + Error::internal_error("failed to fetch resource info", None) + })?; + tx.commit() + .await + .map_err(|_e| Error::internal_error("failed to commit transaction", None))?; + Ok(rows) + } + + async fn inner_get_resources( + user_db: &UserDB, + authed: &ApiAuthed, + workspace_id: &str, + resource_type: &str, + ) -> Result, Error> { + let mut sqlb = SqlBuilder::select_from("resource as o"); + sqlb.fields(&["o.path", "o.description", "o.resource_type"]); + sqlb.and_where("o.workspace_id = ?".bind(&workspace_id)); + sqlb.and_where("o.resource_type = ?".bind(&resource_type)); + let sql = sqlb.sql().map_err(|_e| { + tracing::error!("failed to build sql: {}", _e); + Error::internal_error("failed to build sql", None) + })?; + let mut tx = user_db + .clone() + .begin(authed) + .await + .map_err(|_e| Error::internal_error("failed to begin transaction", None))?; + let rows = sqlx::query_as::<_, ResourceInfo>(&sql) + .fetch_all(&mut *tx) + .await + .map_err(|_e| { + tracing::error!("Failed to fetch resources: {}", _e); + Error::internal_error("failed to fetch resources", None) + })?; + tx.commit() + .await + .map_err(|_e| Error::internal_error("failed to commit transaction", None))?; + + Ok(rows) + } + + async fn inner_get_flows( + user_db: &UserDB, + authed: &ApiAuthed, + workspace_id: &str, + scope_type: &str, + ) -> Result, Error> { + let mut sqlb = SqlBuilder::select_from("flow as o"); + sqlb.fields(&["o.path", "o.summary", "o.description", "o.schema"]); + if scope_type == "favorites" { + sqlb.join("favorite") + .on("favorite.favorite_kind = 'flow' AND favorite.workspace_id = o.workspace_id AND favorite.path = o.path AND favorite.usr = ?" + .bind(&authed.username)); + } + sqlb.and_where("o.workspace_id = ?".bind(&workspace_id)) + .and_where("o.archived = false") + .and_where("o.draft_only IS NOT TRUE") + .order_by("o.edited_at", false) + .limit(100); + let sql = sqlb.sql().map_err(|_e| { + tracing::error!("failed to build sql: {}", _e); + Error::internal_error("failed to build sql", None) + })?; + let mut tx = user_db + .clone() + .begin(authed) + .await + .map_err(|_e| Error::internal_error("failed to begin transaction", None))?; + let rows = sqlx::query_as::<_, FlowInfo>(&sql) + .fetch_all(&mut *tx) + .await + .map_err(|_e| { + tracing::error!("Failed to fetch flows: {}", _e); + Error::internal_error("failed to fetch flows", None) + })?; + tx.commit() + .await + .map_err(|_e| Error::internal_error("failed to commit transaction", None))?; + Ok(rows) + } + + async fn inner_get_scripts( + user_db: &UserDB, + authed: &ApiAuthed, + workspace_id: &str, + scope_type: &str, + ) -> Result, Error> { + let mut sqlb = SqlBuilder::select_from("script as o"); + sqlb.fields(&["o.path", "o.summary", "o.description", "o.schema"]); + if scope_type == "favorites" { + sqlb.join("favorite") + .on("favorite.favorite_kind = 'script' AND favorite.workspace_id = o.workspace_id AND favorite.path = o.path AND favorite.usr = ?" + .bind(&authed.username)); + } + sqlb.and_where("o.workspace_id = ?".bind(&workspace_id)) + .and_where("o.archived = false") + .and_where("o.draft_only IS NOT TRUE") + .order_by("o.created_at", false) + .limit(100); + let sql = sqlb.sql().map_err(|_e| { + tracing::error!("failed to build sql: {}", _e); + Error::internal_error("failed to build sql", None) + })?; + let mut tx = user_db + .clone() + .begin(authed) + .await + .map_err(|_e| Error::internal_error("failed to begin transaction", None))?; + let rows = sqlx::query_as::<_, ScriptInfo>(&sql) + .fetch_all(&mut *tx) + .await + .map_err(|_e| { + tracing::error!("Failed to fetch scripts: {}", _e); + Error::internal_error("failed to fetch scripts", None) + })?; + tx.commit() + .await + .map_err(|_e| Error::internal_error("failed to commit transaction", None))?; + Ok(rows) + } + + async fn transform_schema_for_resources( + schema: &mut Schema, + user_db: &UserDB, + authed: &ApiAuthed, + w_id: &str, + resources_info: &mut HashMap, + ) -> Result<(), Error> { + for (_key, prop) in schema.properties.iter_mut() { + if let Some(format) = &prop.format { + if format.contains("resource") { + let resource_type_key = + format.split("-").last().unwrap_or_default().to_string(); + + if !resources_info.contains_key(&resource_type_key) { + let fetch_result = async { + let resource_type_info_future = Runner::inner_get_resource_type_info( + user_db, + authed, + &w_id, + &resource_type_key, + ); + let resources_data_future = Runner::inner_get_resources( + user_db, + authed, + &w_id, + &resource_type_key, + ); + let (resource_type_info, resources_data) = + try_join!(resource_type_info_future, resources_data_future)?; + Ok::<_, Error>(ResourceCache { + resource_type: resource_type_info, + resources: resources_data, + }) + } + .await; + + match fetch_result { + Ok(cache_data) => { + resources_info.insert(resource_type_key.clone(), cache_data); + } + Err(e) => { + tracing::error!("Failed to fetch resource cache data: {}", e); + return Err(e); + } + } + } + + if let Some(resource_cache) = resources_info.get(&resource_type_key) { + let resources_count = resource_cache.resources.len(); + + prop.r#type = Some("string".to_string()); + prop.description = Some(format!( + "This is a resource named {} with the following description: {}.\nThe path of the resource should be used to specify the resource.\n{}", + resource_cache.resource_type.name, + resource_cache.resource_type.description.as_deref().unwrap_or("No description"), + if resources_count == 0 { + "This resource does not have any available instances, you should create one from your windmill workspace" + } else if resources_count > 1 { + "This resource has multiple available instances, you should precisely select the one you want to use" + } else { + "There is 1 resource available" + } + )); + + if resources_count > 0 { + prop.oneOf = Some( + resource_cache + .resources + .iter() + .map(|resource| { + serde_json::Value::Object(serde_json::Map::from_iter( + [ + ( + "const".to_string(), + serde_json::Value::String(format!( + "$res:{}", + resource.path.clone() + )), + ), + ( + "title".to_string(), + serde_json::Value::String( + resource + .description + .as_deref() + .unwrap_or("No description") + .to_string(), + ), + ), + ] + .into_iter(), + )) + }) + .collect(), + ); + } + } else { + tracing::error!( + "Resource cache entry unexpectedly missing for key: {}", + resource_type_key + ); + } + } + } + } + Ok(()) + } +} + +impl ServerHandler for Runner { + async fn call_tool( + &self, + request: CallToolRequestParam, + context: RequestContext, + ) -> Result { + let parse_args = |args_opt: Option| -> Result { + args_opt.map(Value::Object).ok_or_else(|| { + Error::invalid_params( + "Missing arguments for tool", + Some(request.name.clone().into()), + ) + }) + }; + + let authed = context + .req_extensions + .get::() + .ok_or_else(|| Error::internal_error("ApiAuthed not found", None))?; + let db = context + .req_extensions + .get::() + .ok_or_else(|| Error::internal_error("DB not found", None))?; + let user_db = context + .req_extensions + .get::() + .ok_or_else(|| Error::internal_error("UserDB not found", None))?; + let args = parse_args(request.arguments)?; + + let (tool_type, path) = Runner::reverse_transform(&request.name).unwrap_or_default(); + + // Convert Value to PushArgsOwned + let push_args = if let Value::Object(map) = args.clone() { + let mut args_hash = HashMap::new(); + for (k, v) in map { + args_hash.insert(k, to_raw_value(&v)); + } + windmill_queue::PushArgsOwned { extra: None, args: args_hash } + } else { + windmill_queue::PushArgsOwned::default() + }; + + let w_id = context.workspace_id.clone(); + let script_or_flow_path = StripPath(path); + let run_query = RunJobQuery::default(); + + let result = if tool_type == "script" { + run_wait_result_script_by_path_internal( + db.clone(), + run_query, + script_or_flow_path, + authed.clone(), + user_db.clone(), + w_id.clone(), + push_args, + None, + ) + .await + } else { + run_wait_result_flow_by_path_internal( + db.clone(), + run_query, + script_or_flow_path, + authed.clone(), + user_db.clone(), + push_args, + w_id.clone(), + None, + ) + .await + }; + + match result { + Ok(response) => { + // Extract the response body as bytes, then convert to a string + let body_bytes = to_bytes(response.into_body(), usize::MAX) + .await + .map_err(|e| { + Error::internal_error(format!("Failed to read response body: {}", e), None) + })?; + let body_str = String::from_utf8(body_bytes.to_vec()).map_err(|e| { + Error::internal_error(format!("Failed to decode response body: {}", e), None) + })?; + Ok(CallToolResult::success(vec![Content::text(body_str)])) + } + Err(e) => Err(Error::internal_error( + format!("Failed to run script: {}", e), + None, + )), + } + } + + async fn list_tools( + &self, + _request: Option, + mut _context: RequestContext, + ) -> Result { + let workspace_id = _context.workspace_id.clone(); + let user_db = _context + .req_extensions + .get::() + .ok_or_else(|| Error::internal_error("UserDB not found", None))?; + let authed = _context + .req_extensions + .get::() + .ok_or_else(|| Error::internal_error("ApiAuthed not found", None))?; + let scope = authed + .scopes + .as_ref() + .and_then(|scopes| scopes.iter().find(|scope| scope.starts_with("mcp:"))); + let scope_type = scope.map_or("all", |scope| scope.split(":").last().unwrap_or("all")); + let mut resources_info: HashMap = HashMap::new(); + + let scripts_fn = Runner::inner_get_scripts(user_db, authed, &workspace_id, scope_type); + let flows_fn = Runner::inner_get_flows(user_db, authed, &workspace_id, scope_type); + let (scripts, flows) = try_join!(scripts_fn, flows_fn)?; + + let mut script_tools: Vec = Vec::with_capacity(scripts.len()); + for script in scripts { + let name = Runner::transform_path(&script.path, "script").unwrap_or_default(); + let description = format!( + "This is a script named {} with the following description: {}.", + script.summary.unwrap_or_default(), + script.description.unwrap_or_default() + ); + let mut schema: Schema = script.schema.map_or_else(Schema::default, |v| { + serde_json::from_value(v).unwrap_or_default() + }); + Runner::transform_schema_for_resources( + &mut schema, + user_db, + authed, + &workspace_id, + &mut resources_info, + ) + .await?; + script_tools.push(Tool { + name: Cow::Owned(name), + description: Some(Cow::Owned(description)), + input_schema: { + let value = serde_json::to_value(schema).unwrap_or_default(); + if let serde_json::Value::Object(map) = value { + Arc::new(map) + } else { + Arc::new(serde_json::Map::new()) + } + }, + annotations: None, + }); + } + + let mut flow_tools: Vec = Vec::with_capacity(flows.len()); + for flow in flows { + let name = Runner::transform_path(&flow.path, "flow").unwrap_or_default(); + let description = format!( + "This is a flow named {} with the following description: {}.", + flow.summary.unwrap_or_default(), + flow.description.unwrap_or_default() + ); + let mut schema: Schema = flow.schema.map_or_else(Schema::default, |v| { + serde_json::from_value(v).unwrap_or_default() + }); + Runner::transform_schema_for_resources( + &mut schema, + user_db, + authed, + &workspace_id, + &mut resources_info, + ) + .await?; + flow_tools.push(Tool { + name: Cow::Owned(name), + description: Some(Cow::Owned(description)), + input_schema: { + let value = serde_json::to_value(schema).unwrap_or_default(); + if let serde_json::Value::Object(map) = value { + Arc::new(map) + } else { + Arc::new(serde_json::Map::new()) + } + }, + annotations: None, + }); + } + + let tools = [script_tools, flow_tools].concat(); + Ok(ListToolsResult { tools, next_cursor: None }) + } + + fn get_info(&self) -> ServerInfo { + ServerInfo { + protocol_version: Default::default(), + capabilities: ServerCapabilities::builder() + .enable_tools() + .enable_tool_list_changed() + .build(), + server_info: Implementation::from_build_env(), + instructions: Some("This server provides a runner tool that can run scripts. Use 'get_scripts' to get the list of scripts.".to_string()), + } + } + + async fn initialize( + &self, + _request: InitializeRequestParam, + _context: RequestContext, + ) -> Result { + Ok(self.get_info()) + } + + async fn list_resources( + &self, + _request: Option, + _context: RequestContext, + ) -> Result { + Ok(ListResourcesResult { resources: vec![], next_cursor: None }) + } + + async fn list_prompts( + &self, + _request: Option, + _context: RequestContext, + ) -> Result { + Ok(ListPromptsResult::default()) + } + + async fn list_resource_templates( + &self, + _request: Option, + _context: RequestContext, + ) -> Result { + Ok(ListResourceTemplatesResult::default()) + } +} + +pub fn setup_mcp_server(addr: SocketAddr, path: &str) -> anyhow::Result<(SseServer, Router)> { + let config = SseServerConfig { + bind: addr, + sse_path: "/sse".to_string(), + post_path: "/message".to_string(), + full_message_path: path.to_string(), + ct: CancellationToken::new(), + sse_keep_alive: None, + }; + + Ok(SseServer::new(config)) +} diff --git a/frontend/src/lib/components/TokensTable.svelte b/frontend/src/lib/components/TokensTable.svelte new file mode 100644 index 0000000000..9058a99e1c --- /dev/null +++ b/frontend/src/lib/components/TokensTable.svelte @@ -0,0 +1,295 @@ + + +
+

Tokens

+
+ +
+
+
+ Authenticate to the Windmill API with access tokens. +
+ +
+ {#if newToken} +
+
+ Added token: +
+
+ Make sure to copy your personal access token now. You won't be able to see it again! +
+
+ {/if} + + {#if newMcpToken} +
+

New MCP URL:

+ +

+ Make sure to copy this URL now. You won't be able to see it again! +

+
+ {/if} + + {#if displayCreateToken} +
+

Add a new token

+ {#if scopes != undefined} + {#each scopes as scope} +
+ + +
+ {/each} + {/if} + {#if showMcpMode} + { + mcpCreationMode = e.detail + if (e.detail) { + newTokenLabel = 'MCP token' + newTokenExpiration = undefined + newTokenWorkspace = $workspaceStore + } + }} + checked={mcpCreationMode} + options={{ + right: 'Generate MCP URL', + rightTooltip: + 'Generate a new MCP URL to make your scripts and flows available as tools.' + }} + class="mb-4" + size="xs" + /> + {/if} +
+ {#if mcpCreationMode} +
+ + + + + +
+
+ + +
+ {/if} +
+ + +
+
+ + +
+
+ +
+
+
+ {/if} +
+ +
+ + + prefix + label + expiration + scopes + + + + {#if tokens && tokens.length > 0} + {#each tokens as { token_prefix, expiration, label, scopes }} + + {token_prefix}**** + {label ?? ''} + {displayDate(expiration ?? '')} + {scopes?.join(', ') ?? ''} + + + + + {/each} + {:else if tokens && tokens.length === 0} + + There are no tokens yet + + {:else} + Loading... + {/if} + + +
+ {#if tokens?.length == 100} + + {/if} + {#if tokenPage > 1} + + {/if} +
+
diff --git a/frontend/src/lib/components/UserSettings.svelte b/frontend/src/lib/components/UserSettings.svelte index c02f998331..7ed78507a5 100644 --- a/frontend/src/lib/components/UserSettings.svelte +++ b/frontend/src/lib/components/UserSettings.svelte @@ -5,38 +5,38 @@ stepInputCompletionEnabled, usersWorkspaceStore } from '$lib/stores' - import type { TruncatedToken, NewToken } from '$lib/gen' + import type { TruncatedToken } from '$lib/gen' import { UserService } from '$lib/gen' - import { displayDate, copyToClipboard, getLocalSetting, storeLocalSetting } from '$lib/utils' - import TableCustom from '$lib/components/TableCustom.svelte' + import { getLocalSetting, storeLocalSetting } from '$lib/utils' import { Button } from '$lib/components/common' import Drawer from '$lib/components/common/drawer/Drawer.svelte' import DrawerContent from '$lib/components/common/drawer/DrawerContent.svelte' import { sendUserToast } from '$lib/toast' - import Tooltip from './Tooltip.svelte' import Version from './Version.svelte' - import { Clipboard, Plus } from 'lucide-svelte' import DarkModeToggle from './sidebar/DarkModeToggle.svelte' import Toggle from './Toggle.svelte' import type { Writable } from 'svelte/store' + import TokensTable from './TokensTable.svelte' import { createEventDispatcher } from 'svelte' export let scopes: string[] | undefined = undefined export let newTokenLabel: string | undefined = undefined export let newTokenWorkspace: string | undefined = undefined export let newToken: string | undefined = undefined + export let showMcpMode: boolean = false let newPassword: string | undefined let passwordError: string | undefined let tokens: TruncatedToken[] - let newTokenExpiration: number | undefined - let displayCreateToken = scopes != undefined let login_type = 'none' let drawer: Drawer + let tokenPage = 1 + let openWithMcpMode = false const dispatch = createEventDispatcher() - export function openDrawer() { + export function openDrawer(mcpMode: boolean = false) { + openWithMcpMode = mcpMode loadLoginType() listTokens() drawer?.openDrawer() @@ -68,28 +68,12 @@ login_type = (await UserService.globalWhoami()).login_type } - async function createToken(): Promise { - newToken = undefined - let date: Date | undefined - if (newTokenExpiration) { - date = new Date(new Date().getTime() + newTokenExpiration * 1000) - } - newToken = await UserService.createToken({ - requestBody: { - label: newTokenLabel, - expiration: date?.toISOString(), - scopes, - workspace_id: newTokenWorkspace - } as NewToken - }) - dispatch('tokenCreated', newToken) - listTokens() - displayCreateToken = false - } - - let tokenPage = 1 async function listTokens(): Promise { - tokens = await UserService.listTokens({ excludeEphemeral: true, page: tokenPage, perPage: 100 }) + tokens = await UserService.listTokens({ + excludeEphemeral: true, + page: tokenPage, + perPage: 100 + }) } async function deleteToken(tokenPrefix: string) { @@ -111,6 +95,21 @@ storeLocalSetting(setting, value.toString()) } + function handleNextPage() { + tokenPage += 1 + listTokens() + } + + function handlePreviousPage() { + tokenPage -= 1 + listTokens() + } + + function handleTokenCreated(event: CustomEvent) { + newToken = event.detail + dispatch('tokenCreated', newToken) + } + loadSettings() @@ -221,157 +220,20 @@ {/if} -
-

Tokens

-
- -
-
-
- Authentify to the Windmill API with access tokens. -
- -
-
-
- Added token: -
-
- Make sure to copy your personal access token now. You won’t be able to see it again! -
-
- - -
-

Add a new token

- {#if scopes != undefined} - {#each scopes as scope} -
- - -
- {/each} - {/if} -
-
- - -
-
- - -
-
- -
-
-
-
-
- - - prefix - label - expiration - scopes - - - - {#if tokens && tokens.length > 0} - {#each tokens as { token_prefix, expiration, label, scopes }} - - {token_prefix}**** - {label ?? ''} - {displayDate(expiration ?? '')} - {scopes?.join(', ') ?? ''} - - - {/each} - {:else if tokens && tokens.length === 0} - There are no tokens yet - {:else} - Loading... - {/if} - - -
- {#if tokens?.length == 100} - - {/if} - {#if tokenPage > 1} - - {/if} -
-
+ diff --git a/frontend/src/lib/components/details/EmailTriggerConfigSection.svelte b/frontend/src/lib/components/details/EmailTriggerConfigSection.svelte index e85e4f9de0..eba48136aa 100644 --- a/frontend/src/lib/components/details/EmailTriggerConfigSection.svelte +++ b/frontend/src/lib/components/details/EmailTriggerConfigSection.svelte @@ -89,7 +89,12 @@ placeholder="paste your token here once created to alter examples below" class="!text-xs" /> - diff --git a/frontend/src/routes/(root)/(logged)/user/(user)/workspaces/+page.svelte b/frontend/src/routes/(root)/(logged)/user/(user)/workspaces/+page.svelte index dc5e4a0c19..016e923464 100644 --- a/frontend/src/routes/(root)/(logged)/user/(user)/workspaces/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/user/(user)/workspaces/+page.svelte @@ -28,15 +28,18 @@ let invites: WorkspaceInvite[] = [] let list_all_as_super_admin: boolean = false - let workspaces: { id: string; name: string; username: string; color?: string | null }[] | undefined = undefined + let workspaces: + | { id: string; name: string; username: string; color?: string | null }[] + | undefined = undefined let userSettings: UserSettings let superadminSettings: SuperadminSettings $: rd = $page.url.searchParams.get('rd') - $: if (userSettings && $page.url.hash === USER_SETTINGS_HASH) { - userSettings.openDrawer() + $: if (userSettings && $page.url.hash.startsWith(USER_SETTINGS_HASH)) { + const mcpMode = $page.url.hash.includes('-mcp') + userSettings.openDrawer(mcpMode) } async function loadInvites() { @@ -197,8 +200,8 @@ > {#if workspace.color} {/if} {workspace.id} - {workspace.name} as @@ -295,7 +298,12 @@ Superadmin settings {/if} - @@ -321,4 +329,4 @@

--> - +