From a0e599e3b9c2ee9c103e2c8a6e7161cf86f910ed Mon Sep 17 00:00:00 2001 From: HugoCasa Date: Mon, 27 Jan 2025 22:40:51 +0100 Subject: [PATCH] feat: test trigger connection (#5145) * feat: test trigger connection * fix build * fix build * update ee ref * nit --------- Co-authored-by: Ruben Fiszel --- backend/windmill-api/openapi.yaml | 87 +++++ backend/windmill-api/src/capture.rs | 20 +- backend/windmill-api/src/http_triggers.rs | 16 +- backend/windmill-api/src/kafka_triggers_ee.rs | 5 +- backend/windmill-api/src/lib.rs | 18 +- backend/windmill-api/src/nats_triggers_ee.rs | 8 +- .../windmill-api/src/postgres_triggers/mod.rs | 7 +- .../src/postgres_triggers/trigger.rs | 7 +- .../windmill-api/src/websocket_triggers.rs | 343 ++++++++++-------- .../src/lib/components/ResourceEditor.svelte | 7 +- .../triggers/TestTriggerConnection.svelte | 72 ++++ .../kafka/KafkaTriggersConfigSection.svelte | 24 +- .../nats/NatsTriggersConfigSection.svelte | 19 +- .../WebsocketEditorConfigSection.svelte | 5 + 14 files changed, 422 insertions(+), 216 deletions(-) create mode 100644 frontend/src/lib/components/triggers/TestTriggerConnection.svelte diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index b6b0f05cf9..28abe897b4 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -7943,6 +7943,36 @@ paths: schema: type: string + /w/{workspace}/websocket_triggers/test: + post: + summary: test websocket connection + operationId: testWebsocketConnection + tags: + - websocket_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + requestBody: + description: test websocket connection + required: true + content: + application/json: + schema: + type: object + properties: + url: + type: string + url_runnable_args: + $ref: "#/components/schemas/ScriptArgs" + required: + - url + responses: + "200": + description: successfuly connected to websocket + content: + text/plain: + schema: + type: string + /w/{workspace}/kafka_triggers/create: post: summary: create kafka trigger @@ -8104,6 +8134,34 @@ paths: schema: type: string + /w/{workspace}/kafka_triggers/test: + post: + summary: test kafka connection + operationId: testKafkaConnection + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + requestBody: + description: test kafka connection + required: true + content: + application/json: + schema: + type: object + properties: + connection: + type: object + required: + - connection + responses: + "200": + description: successfuly connected to kafka brokers + content: + text/plain: + schema: + type: string + /w/{workspace}/nats_triggers/create: post: summary: create nats trigger @@ -8267,6 +8325,35 @@ paths: schema: type: string + + /w/{workspace}/nats_triggers/test: + post: + summary: test NATS connection + operationId: testNatsConnection + tags: + - nats_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + requestBody: + description: test nats connection + required: true + content: + application/json: + schema: + type: object + properties: + connection: + type: object + required: + - connection + responses: + "200": + description: successfuly connected to NATS servers + content: + text/plain: + schema: + type: string + /w/{workspace}/postgres_triggers/is_valid_postgres_configuration/{path}: get: summary: check if postgres configuration is set to logical diff --git a/backend/windmill-api/src/capture.rs b/backend/windmill-api/src/capture.rs index 323cb6d67e..22824791a1 100644 --- a/backend/windmill-api/src/capture.rs +++ b/backend/windmill-api/src/capture.rs @@ -35,9 +35,9 @@ use windmill_queue::{PushArgs, PushArgsOwned}; #[cfg(feature = "http_trigger")] use crate::http_triggers::{build_http_trigger_extra, HttpMethod}; #[cfg(all(feature = "enterprise", feature = "kafka"))] -use crate::kafka_triggers_ee::KafkaResourceSecurity; +use crate::kafka_triggers_ee::KafkaTriggerConfigConnection; #[cfg(all(feature = "enterprise", feature = "nats"))] -use crate::nats_triggers_ee::NatsResourceAuth; +use crate::nats_triggers_ee::NatsTriggerConfigConnection; use crate::{ args::WebhookArgs, db::{ApiAuthed, DB}, @@ -110,14 +110,6 @@ struct HttpTriggerConfig { http_method: HttpMethod, } -#[cfg(all(feature = "enterprise", feature = "kafka"))] -#[derive(Serialize, Deserialize)] -#[serde(untagged)] -pub enum KafkaTriggerConfigConnection { - Resource { kafka_resource_path: String }, - Static { brokers: Vec, security: KafkaResourceSecurity }, -} - #[cfg(all(feature = "enterprise", feature = "kafka"))] #[derive(Serialize, Deserialize)] pub struct KafkaTriggerConfig { @@ -127,14 +119,6 @@ pub struct KafkaTriggerConfig { pub group_id: String, } -#[cfg(all(feature = "enterprise", feature = "nats"))] -#[derive(Serialize, Deserialize)] -#[serde(untagged)] -pub enum NatsTriggerConfigConnection { - Resource { nats_resource_path: String }, - Static { servers: Vec, auth: NatsResourceAuth, require_tls: bool }, -} - #[cfg(all(feature = "enterprise", feature = "nats"))] #[derive(Serialize, Deserialize)] pub struct NatsTriggerConfig { diff --git a/backend/windmill-api/src/http_triggers.rs b/backend/windmill-api/src/http_triggers.rs index 00a2cd915c..ef62928f2c 100644 --- a/backend/windmill-api/src/http_triggers.rs +++ b/backend/windmill-api/src/http_triggers.rs @@ -521,7 +521,7 @@ async fn get_http_route_trigger( trigger.email.clone(), &trigger.workspace_id, &db, - Some(username_override.unwrap_or("anonymous".to_string())), + Some(username_override.unwrap_or(format!("http-{}", trigger.path))), ) .await?; @@ -685,12 +685,6 @@ async fn route_job( .await, ); - let label_prefix = Some(format!( - "http-{}-{}-", - method.as_str().to_lowercase(), - trigger.route_path - )); - let run_query = RunJobQuery::default(); if trigger.is_flow { @@ -703,7 +697,7 @@ async fn route_job( StripPath(trigger.script_path.to_owned()), run_query, args, - label_prefix, + None, ) .await .into_response() @@ -716,7 +710,7 @@ async fn route_job( user_db, args, trigger.workspace_id.clone(), - label_prefix, + None, ) .await .into_response() @@ -731,7 +725,7 @@ async fn route_job( StripPath(trigger.script_path.to_owned()), run_query, args, - label_prefix, + None, ) .await .into_response() @@ -744,7 +738,7 @@ async fn route_job( user_db, trigger.workspace_id.clone(), args, - label_prefix, + None, ) .await .into_response() diff --git a/backend/windmill-api/src/kafka_triggers_ee.rs b/backend/windmill-api/src/kafka_triggers_ee.rs index e3de5e6042..d02ea970e3 100644 --- a/backend/windmill-api/src/kafka_triggers_ee.rs +++ b/backend/windmill-api/src/kafka_triggers_ee.rs @@ -9,9 +9,12 @@ pub fn workspaced_service() -> Router { Router::new() } -pub async fn start_kafka_consumers( +pub fn start_kafka_consumers( _db: DB, mut _killpill_rx: tokio::sync::broadcast::Receiver<()>, ) -> () { // implementation is not open source } + +#[derive(Serialize, Deserialize)] +pub enum KafkaTriggerConfigConnection {} diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index 006e52f950..36552f9a57 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -59,8 +59,6 @@ mod auth; mod capture; mod concurrency_groups; mod configs; -#[cfg(feature = "postgres_trigger")] -mod postgres_triggers; mod db; mod drafts; pub mod ee; @@ -75,6 +73,8 @@ mod http_triggers; mod indexer_ee; mod inputs; mod integration; +#[cfg(feature = "postgres_trigger")] +mod postgres_triggers; #[cfg(feature = "enterprise")] mod apps_ee; @@ -98,12 +98,12 @@ mod scripts; mod service_logs; mod settings; mod slack_approvals; -#[cfg(feature = "enterprise")] -mod teams_ee; #[cfg(feature = "smtp")] mod smtp_server_ee; mod static_assets; mod stripe_ee; +#[cfg(feature = "enterprise")] +mod teams_ee; mod tracing_init; mod triggers; mod users; @@ -296,24 +296,24 @@ pub async fn run_server( #[cfg(feature = "websocket")] { let ws_killpill_rx = rx.resubscribe(); - websocket_triggers::start_websockets(db.clone(), ws_killpill_rx).await; + websocket_triggers::start_websockets(db.clone(), ws_killpill_rx); } #[cfg(all(feature = "enterprise", feature = "kafka"))] { let kafka_killpill_rx = rx.resubscribe(); - kafka_triggers_ee::start_kafka_consumers(db.clone(), kafka_killpill_rx).await; + kafka_triggers_ee::start_kafka_consumers(db.clone(), kafka_killpill_rx); } #[cfg(all(feature = "enterprise", feature = "nats"))] { let nats_killpill_rx = rx.resubscribe(); - nats_triggers_ee::start_nats_consumers(db.clone(), nats_killpill_rx).await; + nats_triggers_ee::start_nats_consumers(db.clone(), nats_killpill_rx); } #[cfg(feature = "postgres_trigger")] { let db_killpill_rx = rx.resubscribe(); - postgres_triggers::start_database(db.clone(), db_killpill_rx).await; + postgres_triggers::start_database(db.clone(), db_killpill_rx); } } @@ -454,7 +454,7 @@ pub async fn run_server( ) .route("/slack", post(slack_approvals::slack_app_callback_handler)) .nest("/teams", { - #[cfg(feature = "enterprise")] + #[cfg(feature = "enterprise")] { teams_ee::teams_service() } diff --git a/backend/windmill-api/src/nats_triggers_ee.rs b/backend/windmill-api/src/nats_triggers_ee.rs index fedf7981ff..bfa5fbdd73 100644 --- a/backend/windmill-api/src/nats_triggers_ee.rs +++ b/backend/windmill-api/src/nats_triggers_ee.rs @@ -9,9 +9,9 @@ pub fn workspaced_service() -> Router { Router::new() } -pub async fn start_nats_consumers( - _db: DB, - mut _killpill_rx: tokio::sync::broadcast::Receiver<()>, -) -> () { +pub fn start_nats_consumers(_db: DB, mut _killpill_rx: tokio::sync::broadcast::Receiver<()>) -> () { // implementation is not open source } + +#[derive(Serialize, Deserialize)] +pub enum NatsTriggerConfigConnection {} diff --git a/backend/windmill-api/src/postgres_triggers/mod.rs b/backend/windmill-api/src/postgres_triggers/mod.rs index 9aa50f8fec..0bd87ed262 100644 --- a/backend/windmill-api/src/postgres_triggers/mod.rs +++ b/backend/windmill-api/src/postgres_triggers/mod.rs @@ -110,14 +110,13 @@ async fn run_job( trigger: &PostgresTrigger, ) -> anyhow::Result<()> { let args = PushArgsOwned { args: args.unwrap_or_default(), extra }; - let label_prefix = Some(format!("db-{}-", trigger.path)); let authed = fetch_api_authed( trigger.edited_by.clone(), trigger.email.clone(), &trigger.workspace_id, db, - Some("anonymous".to_string()), + Some(format!("postgres-{}", trigger.path)), ) .await?; @@ -134,7 +133,7 @@ async fn run_job( StripPath(trigger.script_path.to_owned()), run_query, args, - label_prefix, + None, ) .await?; } else { @@ -146,7 +145,7 @@ async fn run_job( StripPath(trigger.script_path.to_owned()), run_query, args, - label_prefix, + None, ) .await?; } diff --git a/backend/windmill-api/src/postgres_triggers/trigger.rs b/backend/windmill-api/src/postgres_triggers/trigger.rs index 09704d0e3b..bd40f41c85 100644 --- a/backend/windmill-api/src/postgres_triggers/trigger.rs +++ b/backend/windmill-api/src/postgres_triggers/trigger.rs @@ -305,7 +305,10 @@ async fn listen_to_transactions( let message = match message { Some(message) => message, None => { - tracing::info!("Stream for postgres trigger {} is empty, leaving....", postgres_trigger.path); + tracing::error!("Stream for postgres trigger {} closed", postgres_trigger.path); + if let None = update_ping(&db, postgres_trigger, Some("Stream closed")).await { + return; + } return; } }; @@ -506,7 +509,7 @@ async fn listen_to_unlistened_database_events( }; } -pub async fn start_database(db: DB, mut killpill_rx: tokio::sync::broadcast::Receiver<()>) { +pub fn start_database(db: DB, mut killpill_rx: tokio::sync::broadcast::Receiver<()>) { tokio::spawn(async move { listen_to_unlistened_database_events(&db, &killpill_rx).await; loop { diff --git a/backend/windmill-api/src/websocket_triggers.rs b/backend/windmill-api/src/websocket_triggers.rs index d55b002598..7e70c602bd 100644 --- a/backend/windmill-api/src/websocket_triggers.rs +++ b/backend/windmill-api/src/websocket_triggers.rs @@ -48,6 +48,7 @@ pub fn workspaced_service() -> Router { .route("/delete/*path", delete(delete_websocket_trigger)) .route("/exists/*path", get(exists_websocket_trigger)) .route("/setenabled/*path", post(set_enabled)) + .route("/test", post(test_websocket_connection)) } #[derive(Deserialize)] @@ -377,6 +378,64 @@ async fn exists_websocket_trigger( Ok(Json(exists)) } +#[derive(Debug, Deserialize)] +struct TestWebsocket { + url: String, + url_runnable_args: Option>, +} + +async fn test_websocket_connection( + authed: ApiAuthed, + Extension(db): Extension, + Path(workspace_id): Path, + Json(test_websocket): Json, +) -> error::Result<()> { + let url = test_websocket.url; + + let connect_f = async { + let connect_url: Cow = if url.starts_with("$") { + if url.starts_with("$flow:") || url.starts_with("$script:") { + let path = url.splitn(2, ':').nth(1).unwrap(); + Cow::Owned( + get_url_from_runnable( + path, + url.starts_with("$flow:"), + &db, + authed, + test_websocket.url_runnable_args.as_ref(), + &workspace_id, + ) + .await?, + ) + } else { + return Err(error::Error::BadConfig(format!( + "Invalid websocket runnable path: {}", + url + ))); + } + } else { + Cow::Borrowed(&url) + }; + + connect_async(connect_url.as_ref()).await.map_err(|err| { + error::Error::BadConfig(format!( + "Error connecting to websocket: {}", + err.to_string() + )) + })?; + + Ok(()) + }; + + tokio::time::timeout(tokio::time::Duration::from_secs(30), connect_f) + .await + .map_err(|_| { + error::Error::BadConfig(format!("Timeout connecting to websocket after 30 seconds")) + })??; + + Ok(()) +} + async fn listen_to_unlistened_websockets( db: &DB, killpill_rx: &tokio::sync::broadcast::Receiver<()>, @@ -419,7 +478,7 @@ async fn listen_to_unlistened_websockets( } } -pub async fn start_websockets(db: DB, mut killpill_rx: tokio::sync::broadcast::Receiver<()>) -> () { +pub fn start_websockets(db: DB, mut killpill_rx: tokio::sync::broadcast::Receiver<()>) -> () { tokio::spawn(async move { listen_to_unlistened_websockets(&db, &killpill_rx).await; loop { @@ -509,7 +568,6 @@ async fn wait_runnable_result( authed: ApiAuthed, db: &DB, workspace_id: &str, - trigger_path: &str, ) -> error::Result { let user_db = UserDB::new(db.clone()); @@ -521,7 +579,6 @@ async fn wait_runnable_result( HashMap::new() }; - let label_prefix = Some(format!("ws-{}-", trigger_path)); let (_, job_id) = if is_flow { run_flow_by_path_inner( authed, @@ -531,7 +588,7 @@ async fn wait_runnable_result( StripPath(path.clone()), RunJobQuery::default(), PushArgsOwned { args, extra: None }, - label_prefix, + None, ) .await? } else { @@ -543,7 +600,7 @@ async fn wait_runnable_result( StripPath(path.clone()), RunJobQuery::default(), PushArgsOwned { args, extra: None }, - label_prefix, + None, ) .await? }; @@ -553,7 +610,9 @@ async fn wait_runnable_result( loop { if start_time.elapsed() > tokio::time::Duration::from_secs(300) { return Err(anyhow::anyhow!( - "Timed out after 5m waiting for runnable {path} (is_flow: {is_flow}) to complete", + "Timed out after 5m waiting for {} {} to complete", + if is_flow { "flow" } else { "script" }, + path ) .into()); } @@ -576,7 +635,8 @@ async fn wait_runnable_result( Ok(Some(r)) => { if !r.success { return Err(anyhow::anyhow!( - "Runnable {path} (is_flow: {is_flow}) failed: {:?}", + "{} {path} failed: {:?}", + if is_flow { "Flow" } else { "Script" }, r.result ) .into()); @@ -590,7 +650,8 @@ async fn wait_runnable_result( } Err(err) => { return Err(anyhow::anyhow!( - "Error fetching job result for runnable {path} (is_flow: {is_flow}): {err}", + "Error fetching job result for {} {path}: {err}", + if is_flow { "flow" } else { "script" }, ) .into()); } @@ -614,25 +675,24 @@ async fn get_url_from_runnable( authed: ApiAuthed, args: Option<&Box>, workspace_id: &str, - trigger_path: &str, ) -> error::Result { - tracing::info!("Running runnable {path} (is_flow: {is_flow}) to get websocket URL",); + tracing::info!( + "Running {} {} to get websocket URL", + if is_flow { "flow" } else { "script" }, + path + ); - let result = wait_runnable_result( - path.to_string(), - is_flow, - args, - authed, - db, - workspace_id, - trigger_path, - ) - .await?; + let result = + wait_runnable_result(path.to_string(), is_flow, args, authed, db, workspace_id).await?; if result.starts_with("\"") && result.ends_with("\"") { Ok(result[1..result.len() - 1].to_string()) } else { - Err(anyhow::anyhow!("Runnable {path} (is_flow: {is_flow}) did not return a string").into()) + Err(error::Error::BadConfig(format!( + "{} {} did not return a string", + if is_flow { "Flow" } else { "Script" }, + path + ))) } } @@ -721,10 +781,9 @@ impl WebsocketTrigger { &path, is_flow, db, - self.fetch_authed(db, Some("url".to_string())).await?, + self.fetch_authed(db).await?, self.url_runnable_args.as_ref().map(|r| &r.0), &self.workspace_id, - &self.path, ) .await } @@ -763,7 +822,9 @@ impl WebsocketTrigger { } InitialMessage::RunnableResult { path, is_flow, args } => { tracing::info!( - "Running runnable {path} (is_flow: {is_flow}) for initial message to websocket {}", + "Running {} {} for initial message to websocket {}", + if is_flow { "flow" } else { "script" }, + path, self.url, ); @@ -771,15 +832,16 @@ impl WebsocketTrigger { path.clone(), is_flow, Some(&args), - self.fetch_authed(db, Some("init".to_string())).await?, + self.fetch_authed(db).await?, db, &self.workspace_id, - &self.path, ) .await?; tracing::info!( - "Sending runnable {path} (is_flow: {is_flow}) result to websocket {}", + "Sending {} {} result to websocket {}", + if is_flow { "flow" } else { "script" }, + path, self.url ); @@ -794,7 +856,11 @@ impl WebsocketTrigger { .await .map_err(to_anyhow) .with_context(|| { - format!("Failed to send runnable {path} (is_flow: {is_flow}) result") + format!( + "Failed to send {} {} result", + if is_flow { "flow" } else { "script" }, + path + ) })?; } } @@ -818,17 +884,13 @@ impl WebsocketTrigger { }; } - async fn fetch_authed( - &self, - db: &DB, - username_override: Option, - ) -> error::Result { + async fn fetch_authed(&self, db: &DB) -> error::Result { fetch_api_authed( self.edited_by.clone(), self.email.clone(), &self.workspace_id, db, - username_override, + Some(format!("ws-{}", self.path)), ) .await } @@ -932,25 +994,20 @@ impl CaptureConfigForWebsocket { &path, is_flow, db, - self.fetch_authed(db, Some("url".to_string())).await?, + self.fetch_authed(db).await?, url_runnable_args.as_ref(), &self.workspace_id, - &self.get_trigger_path(), ) .await } - async fn fetch_authed( - &self, - db: &DB, - username_override: Option, - ) -> error::Result { + async fn fetch_authed(&self, db: &DB) -> error::Result { fetch_api_authed( self.owner.clone(), self.email.clone(), &self.workspace_id, db, - username_override, + Some(format!("ws-{}", self.get_trigger_path())), ) .await } @@ -1017,13 +1074,9 @@ async fn listen_to_websocket( db: DB, mut killpill_rx: tokio::sync::broadcast::Receiver<()>, ) -> () { - if let None = ws.update_ping(&db, Some("Connecting")).await { - return; - } - let url = match &ws { - WebsocketEnum::Trigger(ws_trigger) => &ws_trigger.url, - WebsocketEnum::Capture(capture) => &capture.trigger_config.url, + WebsocketEnum::Trigger(ws_trigger) => ws_trigger.url.clone(), + WebsocketEnum::Capture(capture) => capture.trigger_config.url.clone(), }; let filters: Vec = match &ws { @@ -1035,101 +1088,92 @@ async fn listen_to_websocket( WebsocketEnum::Capture(_) => vec![], }; - loop { - let connect_url: Cow = if url.starts_with("$") { - if url.starts_with("$flow:") || url.starts_with("$script:") { - let path = url.splitn(2, ':').nth(1).unwrap(); - tokio::select! { - biased; - _ = killpill_rx.recv() => { - return; - }, - _ = loop_ping(&db, &ws, Some( - "Waiting on runnable to return websocket URL..." - )) => { - return; - }, + let connect_url: Cow = if url.starts_with("$") { + if url.starts_with("$flow:") || url.starts_with("$script:") { + let path = url.splitn(2, ':').nth(1).unwrap(); + tokio::select! { + biased; + _ = killpill_rx.recv() => { + return; + }, + _ = loop_ping(&db, &ws, Some( + "Waiting on runnable to return websocket URL..." + )) => { + return; + }, - - url_result = ws.get_url_from_runnable(path, url.starts_with("$flow:"), &db) => match url_result { - Ok(url) => Cow::Owned(url), - Err(err) => { - ws.disable_with_error(&db, format!( - "Error getting websocket URL from runnable after 5 tries: {:?}", - err - ), - ) - .await; - return; - } - }, - } - } else { - ws.disable_with_error(&db, format!("Invalid websocket runnable path: {}", url)) - .await; - return; + url_result = ws.get_url_from_runnable(path, url.starts_with("$flow:"), &db) => match url_result { + Ok(url) => Cow::Owned(url), + Err(err) => { + ws.disable_with_error(&db, format!( + "Error getting websocket URL from runnable after 5 tries: {:?}", + err + ), + ) + .await; + return; + } + }, } } else { - Cow::Borrowed(url) - }; + ws.disable_with_error(&db, format!("Invalid websocket runnable path: {}", url)) + .await; + return; + } + } else { + Cow::Borrowed(&url) + }; - tokio::select! { - biased; - _ = killpill_rx.recv() => { - return; - }, - _ = loop_ping(&db, &ws, Some("Connecting...")) => { - return; - }, - connection = connect_async(connect_url.as_ref()) => { - match connection { - Ok((ws_stream, _)) => { - tracing::info!("Listening to websocket {}", url); - if let None = ws.update_ping(&db, None).await { - return; - } - let (writer, mut reader) = ws_stream.split(); - let mut last_ping = tokio::time::Instant::now(); + tokio::select! { + biased; + _ = killpill_rx.recv() => { + return; + }, + _ = loop_ping(&db, &ws, Some("Connecting...")) => { + return; + }, + connection = connect_async(connect_url.as_ref()) => { + match connection { + Ok((ws_stream, _)) => { + tracing::info!("Connected to websocket {}", url); + let (writer, mut reader) = ws_stream.split(); - - // send initial messages - match &ws { - WebsocketEnum::Trigger(ws_trigger) => { - tokio::select! { - biased; - _ = killpill_rx.recv() => { - return; - }, - _ = loop_ping(&db, &ws, Some("Sending initial messages...")) => { - return; - }, - result = ws_trigger.send_initial_messages(writer, &db) => { - if let Err(err) = result { - ws_trigger.disable_with_error(&db, format!("Error sending initial messages: {:?}", err)).await; - return - } else { - tracing::debug!("Initial messages sent successfully to websocket {}", url); - } - } - } - }, - _ => {} - } - - loop { + // send initial messages + match &ws { + WebsocketEnum::Trigger(ws_trigger) => { tokio::select! { biased; _ = killpill_rx.recv() => { return; + }, + _ = loop_ping(&db, &ws, Some("Sending initial messages...")) => { + return; + }, + result = ws_trigger.send_initial_messages(writer, &db) => { + if let Err(err) = result { + ws_trigger.disable_with_error(&db, format!("Error sending initial messages: {:?}", err)).await; + return + } else { + tracing::debug!("Initial messages sent successfully to websocket {}", url); + } } - msg = reader.next() => { - if let Some(msg) = msg { - if last_ping.elapsed() > tokio::time::Duration::from_secs(5) { - if let None = ws.update_ping(&db, None).await { - return; - } - last_ping = tokio::time::Instant::now(); - } + } + }, + _ => {} + } + + loop { + tokio::select! { + biased; + _ = killpill_rx.recv() => { + return; + }, + _ = loop_ping(&db, &ws, None) => { + return; + }, + _ = async { + loop { + if let Some(msg) = reader.next().await { match msg { Ok(msg) => { match msg { @@ -1170,7 +1214,7 @@ async fn listen_to_websocket( }, WebsocketEnum::Capture(capture) => { capture.handle(&db, args).await; - } + }, } } }, @@ -1188,26 +1232,19 @@ async fn listen_to_websocket( if let None = ws.update_ping(&db, Some("Websocket closed")).await { return; } - tokio::time::sleep(tokio::time::Duration::from_secs(5)).await; - break; - } - }, - _ = tokio::time::sleep(tokio::time::Duration::from_secs(5)) => { - tracing::debug!("Sending ping to websocket {}", url); - if let None = ws.update_ping(&db, None).await { return; } - last_ping = tokio::time::Instant::now(); - }, + } + } => { + return; } } } - Err(err) => { - tracing::error!("Error connecting to websocket {}: {:?}", url, err); - if let None = ws.update_ping(&db, Some(err.to_string().as_str())).await { - return; - } - tokio::time::sleep(tokio::time::Duration::from_secs(5)).await; + } + Err(err) => { + tracing::error!("Error connecting to websocket {}: {:?}", url, err); + if let None = ws.update_ping(&db, Some(err.to_string().as_str())).await { + return; } } } @@ -1216,14 +1253,12 @@ async fn listen_to_websocket( } async fn run_job(db: &DB, trigger: &WebsocketTrigger, args: PushArgsOwned) -> anyhow::Result<()> { - let label_prefix = Some(format!("ws-{}-", trigger.path)); - let authed = fetch_api_authed( trigger.edited_by.clone(), trigger.email.clone(), &trigger.workspace_id, db, - Some("anonymous".to_string()), + Some(format!("ws-{}", trigger.path)), ) .await?; @@ -1240,7 +1275,7 @@ async fn run_job(db: &DB, trigger: &WebsocketTrigger, args: PushArgsOwned) -> an StripPath(trigger.script_path.to_owned()), run_query, args, - label_prefix, + None, ) .await?; } else { @@ -1252,7 +1287,7 @@ async fn run_job(db: &DB, trigger: &WebsocketTrigger, args: PushArgsOwned) -> an StripPath(trigger.script_path.to_owned()), run_query, args, - label_prefix, + None, ) .await?; } diff --git a/frontend/src/lib/components/ResourceEditor.svelte b/frontend/src/lib/components/ResourceEditor.svelte index d87268aaf8..8f882baa59 100644 --- a/frontend/src/lib/components/ResourceEditor.svelte +++ b/frontend/src/lib/components/ResourceEditor.svelte @@ -17,6 +17,7 @@ import Markdown from 'svelte-exmarkdown' import autosize from '$lib/autosize' import GfmMarkdown from './GfmMarkdown.svelte' + import TestTriggerConnection from './triggers/TestTriggerConnection.svelte' export let canSave = true export let resource_type: string | undefined = undefined @@ -223,7 +224,11 @@ {/if}
- + {#if resourceToEdit?.resource_type === 'nats' || resourceToEdit?.resource_type === 'kafka'} + + {:else} + + {/if} switchTab(e.detail)} options={{ diff --git a/frontend/src/lib/components/triggers/TestTriggerConnection.svelte b/frontend/src/lib/components/triggers/TestTriggerConnection.svelte new file mode 100644 index 0000000000..8c56f3627b --- /dev/null +++ b/frontend/src/lib/components/triggers/TestTriggerConnection.svelte @@ -0,0 +1,72 @@ + + +
+ +
diff --git a/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte b/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte index 947ce0cfdd..6084de018a 100644 --- a/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte +++ b/frontend/src/lib/components/triggers/kafka/KafkaTriggersConfigSection.svelte @@ -9,6 +9,7 @@ import CaptureSection, { type CaptureInfo } from '../CaptureSection.svelte' import CaptureTable from '../CaptureTable.svelte' import { workspaceStore } from '$lib/stores' + import TestTriggerConnection from '../TestTriggerConnection.svelte' export let path: string export let defaultValues: Record | undefined = undefined export let headless: boolean = false @@ -135,16 +136,19 @@ required: ['topics', 'group_id'] } - let connectionValid = false + let isStaticConnectionValid = false let otherArgsValid = false - $: isValid = - (selected === 'resource' + $: isConnectionValid = + selected === 'resource' ? !!args.kafka_resource_path - : connectionValid && + : isStaticConnectionValid && args.brokers && args.brokers.length > 0 && - args.brokers.every((b) => b.length > 0)) && + args.brokers.every((b) => b.length > 0) + + $: isValid = + isConnectionValid && otherArgsValid && args.topics && args.topics.length > 0 && @@ -213,10 +217,18 @@ {/if} + {#if isConnectionValid} + + {/if}
diff --git a/frontend/src/lib/components/triggers/nats/NatsTriggersConfigSection.svelte b/frontend/src/lib/components/triggers/nats/NatsTriggersConfigSection.svelte index f3ae57c843..7d65ac60c8 100644 --- a/frontend/src/lib/components/triggers/nats/NatsTriggersConfigSection.svelte +++ b/frontend/src/lib/components/triggers/nats/NatsTriggersConfigSection.svelte @@ -9,6 +9,7 @@ import SchemaForm from '$lib/components/SchemaForm.svelte' import CaptureSection, { type CaptureInfo } from '../CaptureSection.svelte' import CaptureTable from '../CaptureTable.svelte' + import TestTriggerConnection from '../TestTriggerConnection.svelte' export let defaultValues: Record | undefined = undefined export let headless: boolean = false export let args: Record = {} @@ -172,7 +173,7 @@ required: ['subjects', 'use_jetstream', 'stream_name', 'consumer_name'] } - let connectionValid = false + let isStaticConnectionValid = false let otherArgsValid = false let globalError = '' @@ -181,15 +182,18 @@ ? 'Only one subject is supported if not using JetStream.' : '' - $: isValid = - (selected === 'resource' + $: isConnectionValid = + selected === 'resource' ? !!args.nats_resource_path - : connectionValid && + : isStaticConnectionValid && args.servers && args.servers.length > 0 && args.servers.every((b) => b.length > 0) && args.require_tls !== undefined && - args.require_tls !== null) && + args.require_tls !== null + + $: isValid = + isConnectionValid && otherArgsValid && args.subjects && args.subjects.length > 0 && @@ -264,10 +268,13 @@ {/if} + {#if isConnectionValid} + + {/if}
diff --git a/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte b/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte index 3db228cf6b..aac583c09d 100644 --- a/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte +++ b/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte @@ -11,6 +11,7 @@ import type { Schema } from '$lib/common' import { FlowService, ScriptService, type Flow, type Script } from '$lib/gen' import { workspaceStore } from '$lib/stores' + import TestTriggerConnection from '../TestTriggerConnection.svelte' export let url: string | undefined export let url_runnable_args: Record | undefined @@ -182,5 +183,9 @@ {/if} + + {#if isValid} + + {/if}