From 487c273bfbc5bf1731947c922f32cd026cca47f9 Mon Sep 17 00:00:00 2001 From: HugoCasa Date: Wed, 29 Jan 2025 21:17:36 +0100 Subject: [PATCH] feat: websocket trigger allow returning messages (#5168) * feat: websocket trigger allow returning messages * handle early return --- ...f0a2db057141619e9818ddb904d7c242efec3.json | 40 -- ...200a07c5a4ad3a50051b73f382385b89466e.json} | 5 +- ...dfdfc49981568e8c496fb6a163c499c3e4ad1.json | 23 + ...4837_websocket_can_return_message.down.sql | 1 + ...094837_websocket_can_return_message.up.sql | 1 + backend/windmill-api/openapi.yaml | 13 +- backend/windmill-api/src/jobs.rs | 196 ++++--- .../windmill-api/src/websocket_triggers.rs | 542 ++++++++++-------- .../details/DetailPageTriggerPanel.svelte | 2 +- .../components/sidebar/SidebarContent.svelte | 2 +- .../triggers/TestTriggerConnection.svelte | 3 +- .../components/triggers/TriggersEditor.svelte | 2 +- .../triggers/TriggersEditorSection.svelte | 4 +- .../kafka/KafkaTriggerEditorInner.svelte | 8 +- .../triggers/kafka/KafkaTriggersPanel.svelte | 4 +- .../WebsocketEditorConfigSection.svelte | 2 +- .../WebsocketTriggerEditorInner.svelte | 61 +- .../websocket/WebsocketTriggersPanel.svelte | 10 +- .../(logged)/kafka_triggers/+page.svelte | 417 +++++++------- .../(logged)/websocket_triggers/+page.js | 2 +- .../(logged)/websocket_triggers/+page.svelte | 16 +- 21 files changed, 740 insertions(+), 614 deletions(-) delete mode 100644 backend/.sqlx/query-1ef48cc430870ab6c062b046bdbf0a2db057141619e9818ddb904d7c242efec3.json rename backend/.sqlx/{query-561b7935d687f6b9f3d6488f8489f55e6737ffa2d3716f803d32f8b68cc1e915.json => query-63b42286804a3f3977235935d7ac200a07c5a4ad3a50051b73f382385b89466e.json} (62%) create mode 100644 backend/.sqlx/query-fc243af1bc70f04e28c006364d6dfdfc49981568e8c496fb6a163c499c3e4ad1.json create mode 100644 backend/migrations/20250129094837_websocket_can_return_message.down.sql create mode 100644 backend/migrations/20250129094837_websocket_can_return_message.up.sql diff --git a/backend/.sqlx/query-1ef48cc430870ab6c062b046bdbf0a2db057141619e9818ddb904d7c242efec3.json b/backend/.sqlx/query-1ef48cc430870ab6c062b046bdbf0a2db057141619e9818ddb904d7c242efec3.json deleted file mode 100644 index 7fc5a4c430..0000000000 --- a/backend/.sqlx/query-1ef48cc430870ab6c062b046bdbf0a2db057141619e9818ddb904d7c242efec3.json +++ /dev/null @@ -1,40 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT \n EXISTS(SELECT 1 FROM websocket_trigger WHERE workspace_id = $1) as \"websocket_used!\", \n EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1) as \"http_routes_used!\",\n EXISTS(SELECT 1 FROM kafka_trigger WHERE workspace_id = $1) as \"kafka_used!\",\n EXISTS(SELECT 1 FROM nats_trigger WHERE workspace_id = $1) as \"nats_used!\"", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "websocket_used!", - "type_info": "Bool" - }, - { - "ordinal": 1, - "name": "http_routes_used!", - "type_info": "Bool" - }, - { - "ordinal": 2, - "name": "kafka_used!", - "type_info": "Bool" - }, - { - "ordinal": 3, - "name": "nats_used!", - "type_info": "Bool" - } - ], - "parameters": { - "Left": [ - "Text" - ] - }, - "nullable": [ - null, - null, - null, - null - ] - }, - "hash": "1ef48cc430870ab6c062b046bdbf0a2db057141619e9818ddb904d7c242efec3" -} diff --git a/backend/.sqlx/query-561b7935d687f6b9f3d6488f8489f55e6737ffa2d3716f803d32f8b68cc1e915.json b/backend/.sqlx/query-63b42286804a3f3977235935d7ac200a07c5a4ad3a50051b73f382385b89466e.json similarity index 62% rename from backend/.sqlx/query-561b7935d687f6b9f3d6488f8489f55e6737ffa2d3716f803d32f8b68cc1e915.json rename to backend/.sqlx/query-63b42286804a3f3977235935d7ac200a07c5a4ad3a50051b73f382385b89466e.json index a3e2f3d338..3938de3ddf 100644 --- a/backend/.sqlx/query-561b7935d687f6b9f3d6488f8489f55e6737ffa2d3716f803d32f8b68cc1e915.json +++ b/backend/.sqlx/query-63b42286804a3f3977235935d7ac200a07c5a4ad3a50051b73f382385b89466e.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE websocket_trigger SET url = $1, script_path = $2, path = $3, is_flow = $4, filters = $5, initial_messages = $6, url_runnable_args = $7, edited_by = $8, email = $9, edited_at = now(), server_id = NULL, error = NULL\n WHERE workspace_id = $10 AND path = $11", + "query": "UPDATE websocket_trigger SET url = $1, script_path = $2, path = $3, is_flow = $4, filters = $5, initial_messages = $6, url_runnable_args = $7, edited_by = $8, email = $9, can_return_message = $10, edited_at = now(), server_id = NULL, error = NULL\n WHERE workspace_id = $11 AND path = $12", "describe": { "columns": [], "parameters": { @@ -14,11 +14,12 @@ "Jsonb", "Varchar", "Varchar", + "Bool", "Text", "Text" ] }, "nullable": [] }, - "hash": "561b7935d687f6b9f3d6488f8489f55e6737ffa2d3716f803d32f8b68cc1e915" + "hash": "63b42286804a3f3977235935d7ac200a07c5a4ad3a50051b73f382385b89466e" } diff --git a/backend/.sqlx/query-fc243af1bc70f04e28c006364d6dfdfc49981568e8c496fb6a163c499c3e4ad1.json b/backend/.sqlx/query-fc243af1bc70f04e28c006364d6dfdfc49981568e8c496fb6a163c499c3e4ad1.json new file mode 100644 index 0000000000..57ae674eb4 --- /dev/null +++ b/backend/.sqlx/query-fc243af1bc70f04e28c006364d6dfdfc49981568e8c496fb6a163c499c3e4ad1.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT flow_version.value->>'early_return' as early_return\n FROM flow \n LEFT JOIN flow_version\n ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]\n WHERE flow.path = $1 and flow.workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "early_return", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "fc243af1bc70f04e28c006364d6dfdfc49981568e8c496fb6a163c499c3e4ad1" +} diff --git a/backend/migrations/20250129094837_websocket_can_return_message.down.sql b/backend/migrations/20250129094837_websocket_can_return_message.down.sql new file mode 100644 index 0000000000..02cd2ce8c2 --- /dev/null +++ b/backend/migrations/20250129094837_websocket_can_return_message.down.sql @@ -0,0 +1 @@ +ALTER TABLE websocket_trigger DROP COLUMN can_return_message; \ No newline at end of file diff --git a/backend/migrations/20250129094837_websocket_can_return_message.up.sql b/backend/migrations/20250129094837_websocket_can_return_message.up.sql new file mode 100644 index 0000000000..3f2a0b17a0 --- /dev/null +++ b/backend/migrations/20250129094837_websocket_can_return_message.up.sql @@ -0,0 +1 @@ +ALTER TABLE websocket_trigger ADD COLUMN can_return_message BOOLEAN NOT NULL DEFAULT FALSE; \ No newline at end of file diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 24e3997824..d3471c28dc 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -7963,8 +7963,11 @@ paths: type: string url_runnable_args: $ref: "#/components/schemas/ScriptArgs" + can_return_message: + type: boolean required: - url + - can_return_message responses: "200": description: successfuly connected to websocket @@ -13139,6 +13142,8 @@ components: $ref: "#/components/schemas/WebsocketTriggerInitialMessage" url_runnable_args: $ref: "#/components/schemas/ScriptArgs" + can_return_message: + type: boolean required: - path @@ -13152,6 +13157,7 @@ components: - workspace_id - enabled - filters + - can_return_message NewWebsocketTrigger: type: object @@ -13183,6 +13189,8 @@ components: $ref: "#/components/schemas/WebsocketTriggerInitialMessage" url_runnable_args: $ref: "#/components/schemas/ScriptArgs" + can_return_message: + type: boolean required: - path @@ -13190,6 +13198,7 @@ components: - url - is_flow - filters + - can_return_message EditWebsocketTrigger: type: object @@ -13219,6 +13228,8 @@ components: $ref: "#/components/schemas/WebsocketTriggerInitialMessage" url_runnable_args: $ref: "#/components/schemas/ScriptArgs" + can_return_message: + type: boolean required: - path @@ -13226,7 +13237,7 @@ components: - url - is_flow - filters - + - can_return_message WebsocketTriggerInitialMessage: anyOf: - type: object diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index b5369c0902..f5a6c78df7 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -3527,13 +3527,14 @@ pub struct WindmillCompositeResult { windmill_headers: Option>, result: Option>, } -pub async fn run_wait_result( + +pub async fn run_wait_result_internal( db: &DB, uuid: Uuid, w_id: String, node_id_for_empty_return: Option, username: &str, -) -> error::Result { +) -> error::Result<(Box, bool)> { let mut result = None; let mut success = false; let timeout = TIMEOUT_WAIT_RESULT.read().await.clone().unwrap_or(600); @@ -3605,104 +3606,113 @@ pub async fn run_wait_result( }; tokio::time::sleep(core::time::Duration::from_millis(delay)).await; } + if let Some(result) = result { g.done = true; - - let composite_result = serde_json::from_str::(result.get()); - match composite_result { - Ok(WindmillCompositeResult { - windmill_status_code, - windmill_content_type, - windmill_headers, - result: result_value, - }) => { - if windmill_content_type.is_none() - && windmill_status_code.is_none() - && windmill_headers.is_none() - { - return Ok(( - if success { - StatusCode::OK - } else { - StatusCode::INTERNAL_SERVER_ERROR - }, - Json(result), - ) - .into_response()); - } - - let status_code_or_default = windmill_status_code - .map(|val| match StatusCode::from_u16(val) { - Ok(sc) => Ok(sc), - Err(_) => Err(Error::ExecutionErr("Invalid status code".to_string())), - }) - .unwrap_or_else(|| { - if !success { - Ok(StatusCode::INTERNAL_SERVER_ERROR) - } else if result_value.is_some() { - Ok(StatusCode::OK) - } else { - Ok(StatusCode::NO_CONTENT) - } - })?; - - let mut headers = HeaderMap::new(); - - if let Some(windmill_headers) = windmill_headers { - for (k, v) in windmill_headers { - let k = HeaderName::from_str(k.as_str()).map_err(|err| { - Error::InternalErr(format!("Invalid header name {k}: {err}")) - })?; - let v = HeaderValue::from_str(v.as_str()).map_err(|err| { - Error::InternalErr(format!("Invalid header value {v}: {err}")) - })?; - headers.insert(k, v); - } - } - - if let Some(content_type) = windmill_content_type { - let serialized_json_result = result_value - .map(|val| val.get().to_owned()) - .unwrap_or_else(String::new); - // if the `result` was just a single string, the below removes the surrounding quotes by parsing it as a string. - // it falls back to the original serialized JSON if it doesn't work. - let serialized_result = - serde_json::from_str::(serialized_json_result.as_str()) - .ok() - .unwrap_or(serialized_json_result); - headers.insert( - http::header::CONTENT_TYPE, - HeaderValue::from_str(content_type.as_str()).map_err(|err| { - Error::InternalErr(format!( - "Invalid content type {content_type}: {err}" - )) - })?, - ); - return Ok((status_code_or_default, headers, serialized_result).into_response()); - } - if let Some(result_value) = result_value { - return Ok( - (status_code_or_default, headers, Json(result_value)).into_response() - ); - } else { - Ok((status_code_or_default, headers).into_response()) - } - } - _ => Ok(( - if success { - StatusCode::OK - } else { - StatusCode::INTERNAL_SERVER_ERROR - }, - Json(result), - ) - .into_response()), - } + Ok((result, success)) } else { Err(Error::ExecutionErr(format!("timeout after {}s", timeout))) } } +pub async fn run_wait_result( + db: &DB, + uuid: Uuid, + w_id: String, + node_id_for_empty_return: Option, + username: &str, +) -> error::Result { + let (result, success) = + run_wait_result_internal(db, uuid, w_id, node_id_for_empty_return, username).await?; + + let composite_result = serde_json::from_str::(result.get()); + match composite_result { + Ok(WindmillCompositeResult { + windmill_status_code, + windmill_content_type, + windmill_headers, + result: result_value, + }) => { + if windmill_content_type.is_none() + && windmill_status_code.is_none() + && windmill_headers.is_none() + { + return Ok(( + if success { + StatusCode::OK + } else { + StatusCode::INTERNAL_SERVER_ERROR + }, + Json(result), + ) + .into_response()); + } + + let status_code_or_default = windmill_status_code + .map(|val| match StatusCode::from_u16(val) { + Ok(sc) => Ok(sc), + Err(_) => Err(Error::ExecutionErr("Invalid status code".to_string())), + }) + .unwrap_or_else(|| { + if !success { + Ok(StatusCode::INTERNAL_SERVER_ERROR) + } else if result_value.is_some() { + Ok(StatusCode::OK) + } else { + Ok(StatusCode::NO_CONTENT) + } + })?; + + let mut headers = HeaderMap::new(); + + if let Some(windmill_headers) = windmill_headers { + for (k, v) in windmill_headers { + let k = HeaderName::from_str(k.as_str()).map_err(|err| { + Error::InternalErr(format!("Invalid header name {k}: {err}")) + })?; + let v = HeaderValue::from_str(v.as_str()).map_err(|err| { + Error::InternalErr(format!("Invalid header value {v}: {err}")) + })?; + headers.insert(k, v); + } + } + + if let Some(content_type) = windmill_content_type { + let serialized_json_result = result_value + .map(|val| val.get().to_owned()) + .unwrap_or_else(String::new); + // if the `result` was just a single string, the below removes the surrounding quotes by parsing it as a string. + // it falls back to the original serialized JSON if it doesn't work. + let serialized_result = + serde_json::from_str::(serialized_json_result.as_str()) + .ok() + .unwrap_or(serialized_json_result); + headers.insert( + http::header::CONTENT_TYPE, + HeaderValue::from_str(content_type.as_str()).map_err(|err| { + Error::InternalErr(format!("Invalid content type {content_type}: {err}")) + })?, + ); + return Ok((status_code_or_default, headers, serialized_result).into_response()); + } + if let Some(result_value) = result_value { + return Ok((status_code_or_default, headers, Json(result_value)).into_response()); + } else { + Ok((status_code_or_default, headers).into_response()) + } + } + _ => Ok(( + if success { + StatusCode::OK + } else { + StatusCode::INTERNAL_SERVER_ERROR + }, + Json(result), + ) + .into_response()), + } +} + async fn delete_job_metadata_after_use(db: &DB, job_uuid: Uuid) -> Result<(), Error> { sqlx::query!( "UPDATE completed_job diff --git a/backend/windmill-api/src/websocket_triggers.rs b/backend/windmill-api/src/websocket_triggers.rs index 7e70c602bd..f90bb8def0 100644 --- a/backend/windmill-api/src/websocket_triggers.rs +++ b/backend/windmill-api/src/websocket_triggers.rs @@ -33,7 +33,9 @@ use windmill_queue::PushArgsOwned; use crate::{ capture::{insert_capture_payload, TriggerKind, WebsocketTriggerConfig}, db::{ApiAuthed, DB}, - jobs::{run_flow_by_path_inner, run_script_by_path_inner, RunJobQuery}, + jobs::{ + run_flow_by_path_inner, run_script_by_path_inner, run_wait_result_internal, RunJobQuery, + }, users::fetch_api_authed, }; @@ -61,6 +63,7 @@ struct NewWebsocketTrigger { filters: Vec>, initial_messages: Option>>, url_runnable_args: Option>, + can_return_message: bool, } #[derive(Deserialize)] @@ -101,6 +104,7 @@ pub struct WebsocketTrigger { filters: Vec>>, initial_messages: Option>>>, url_runnable_args: Option>>, + can_return_message: bool, } #[derive(Deserialize)] @@ -112,6 +116,7 @@ struct EditWebsocketTrigger { filters: Vec>, initial_messages: Option>>, url_runnable_args: Option>, + can_return_message: bool, } #[derive(Deserialize)] @@ -189,7 +194,7 @@ async fn create_websocket_trigger( ) -> error::Result<(StatusCode, String)> { if *CLOUD_HOSTED { return Err(error::Error::BadRequest( - "Websocket triggers are not supported on multi-tenant cloud, use dedicated cloud or self-host".to_string(), + "WebSocket triggers are not supported on multi-tenant cloud, use dedicated cloud or self-host".to_string(), )); } @@ -203,7 +208,7 @@ async fn create_websocket_trigger( .map(SqlxJson) .collect_vec(); sqlx::query_as::<_, WebsocketTrigger>( - "INSERT INTO websocket_trigger (workspace_id, path, url, script_path, is_flow, enabled, filters, initial_messages, url_runnable_args, edited_by, email, edited_at) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, now()) RETURNING *", + "INSERT INTO websocket_trigger (workspace_id, path, url, script_path, is_flow, enabled, filters, initial_messages, url_runnable_args, edited_by, can_return_message, email, edited_at) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, now()) RETURNING *", ) .bind(&w_id) .bind(&ct.path) @@ -215,6 +220,7 @@ async fn create_websocket_trigger( .bind(initial_messages.as_slice()) .bind(ct.url_runnable_args.map(SqlxJson)) .bind(&authed.username) + .bind(ct.can_return_message) .bind(&authed.email) .fetch_one(&mut *tx).await?; @@ -253,8 +259,8 @@ async fn update_websocket_trigger( // important to update server_id to NULL to stop current websocket listener sqlx::query!( - "UPDATE websocket_trigger SET url = $1, script_path = $2, path = $3, is_flow = $4, filters = $5, initial_messages = $6, url_runnable_args = $7, edited_by = $8, email = $9, edited_at = now(), server_id = NULL, error = NULL - WHERE workspace_id = $10 AND path = $11", + "UPDATE websocket_trigger SET url = $1, script_path = $2, path = $3, is_flow = $4, filters = $5, initial_messages = $6, url_runnable_args = $7, edited_by = $8, email = $9, can_return_message = $10, edited_at = now(), server_id = NULL, error = NULL + WHERE workspace_id = $11 AND path = $12", ct.url, ct.script_path, ct.path, @@ -264,6 +270,7 @@ async fn update_websocket_trigger( ct.url_runnable_args.map(SqlxJson) as Option>>, &authed.username, &authed.email, + ct.can_return_message, w_id, path, ) @@ -310,7 +317,7 @@ pub async fn set_enabled( w_id, ).fetch_optional(&mut *tx).await?; - not_found_if_none(one_o.flatten(), "Websocket trigger", path)?; + not_found_if_none(one_o.flatten(), "WebSocket trigger", path)?; audit_log( &mut *tx, @@ -326,7 +333,7 @@ pub async fn set_enabled( tx.commit().await?; Ok(format!( - "succesfully updated websocket trigger at path {} to status {}", + "succesfully updated WebSocket trigger at path {} to status {}", path, payload.enabled )) } @@ -359,7 +366,7 @@ async fn delete_websocket_trigger( tx.commit().await?; - Ok(format!("Websocket trigger {path} deleted")) + Ok(format!("WebSocket trigger {path} deleted")) } async fn exists_websocket_trigger( @@ -409,7 +416,7 @@ async fn test_websocket_connection( ) } else { return Err(error::Error::BadConfig(format!( - "Invalid websocket runnable path: {}", + "Invalid WebSocket runnable path: {}", url ))); } @@ -419,7 +426,7 @@ async fn test_websocket_connection( connect_async(connect_url.as_ref()).await.map_err(|err| { error::Error::BadConfig(format!( - "Error connecting to websocket: {}", + "Error connecting to WebSocket: {}", err.to_string() )) })?; @@ -430,7 +437,7 @@ async fn test_websocket_connection( 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")) + error::Error::BadConfig(format!("Timeout connecting to WebSocket after 30 seconds")) })??; Ok(()) @@ -455,7 +462,7 @@ async fn listen_to_unlistened_websockets( } } Err(err) => { - tracing::error!("Error fetching websocket triggers: {:?}", err); + tracing::error!("Error fetching WebSocket triggers: {:?}", err); } }; @@ -473,7 +480,7 @@ async fn listen_to_unlistened_websockets( } } Err(err) => { - tracing::error!("Error fetching capture websocket triggers: {:?}", err); + tracing::error!("Error fetching capture WebSocket triggers: {:?}", err); } } } @@ -561,16 +568,9 @@ where deserializer.deserialize_map(SupersetVisitor { key, value_to_check }) } -async fn wait_runnable_result( - path: String, - is_flow: bool, +fn raw_value_to_args_hashmap( args: Option<&Box>, - authed: ApiAuthed, - db: &DB, - workspace_id: &str, -) -> error::Result { - let user_db = UserDB::new(db.clone()); - +) -> error::Result>> { let args = if let Some(args) = args { serde_json::from_str::>>>(args.get()) .map_err(|e| error::Error::BadRequest(format!("invalid json: {}", e)))? @@ -578,84 +578,89 @@ async fn wait_runnable_result( } else { HashMap::new() }; + Ok(args) +} - let (_, job_id) = if is_flow { - run_flow_by_path_inner( +async fn wait_runnable_result( + path: String, + is_flow: bool, + args: PushArgsOwned, + authed: ApiAuthed, + db: &DB, + workspace_id: &str, +) -> error::Result { + let user_db = UserDB::new(db.clone()); + + let username = authed.display_username().to_owned(); + + let (job_id, early_return) = if is_flow { + let (_, job_id) = run_flow_by_path_inner( authed, db.clone(), user_db, workspace_id.to_string(), StripPath(path.clone()), RunJobQuery::default(), - PushArgsOwned { args, extra: None }, + args, None, ) + .await?; + + let early_return = sqlx::query_scalar!( + r#"SELECT flow_version.value->>'early_return' as early_return + FROM flow + LEFT JOIN flow_version + ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.path = $1 and flow.workspace_id = $2"#, + path, + workspace_id, + ) + .fetch_optional(db) .await? + .flatten(); + + (job_id, early_return) } else { - run_script_by_path_inner( + let (_, job_id) = run_script_by_path_inner( authed, db.clone(), user_db, workspace_id.to_string(), StripPath(path.clone()), RunJobQuery::default(), - PushArgsOwned { args, extra: None }, + args, None, ) - .await? + .await?; + + (job_id, None) }; - let start_time = tokio::time::Instant::now(); - - loop { - if start_time.elapsed() > tokio::time::Duration::from_secs(300) { - return Err(anyhow::anyhow!( - "Timed out after 5m waiting for {} {} to complete", - if is_flow { "flow" } else { "script" }, - path - ) - .into()); - } - - #[derive(sqlx::FromRow)] - struct RawResult { - result: Option>>, - success: bool, - } - - let result = sqlx::query_as::<_, RawResult>( - "SELECT result, success FROM completed_job WHERE id = $1 AND workspace_id = $2", + let (result, success) = run_wait_result_internal( + db, + Uuid::parse_str(&job_id).unwrap(), + workspace_id.to_string(), + early_return, + &username, + ) + .await + .with_context(|| { + format!( + "Error fetching job result for {} {}", + if is_flow { "flow" } else { "script" }, + path ) - .bind(Uuid::parse_str(&job_id).unwrap()) - .bind(workspace_id) - .fetch_optional(db) - .await; + })?; - match result { - Ok(Some(r)) => { - if !r.success { - return Err(anyhow::anyhow!( - "{} {path} failed: {:?}", - if is_flow { "Flow" } else { "Script" }, - r.result - ) - .into()); - } else { - return Ok(r.result.map(|r| r.get().to_owned()).unwrap_or_default()); - } - } - Ok(None) => { - // not yet done, wait for 5s and check again - tokio::time::sleep(tokio::time::Duration::from_secs(5)).await; - } - Err(err) => { - return Err(anyhow::anyhow!( - "Error fetching job result for {} {path}: {err}", - if is_flow { "flow" } else { "script" }, - ) - .into()); - } - } + if !success { + Err(anyhow::anyhow!( + "{} {path} failed: {:?}", + if is_flow { "Flow" } else { "Script" }, + result + ) + .into()) + } else { + Ok(result.get().to_owned()) } } @@ -677,23 +682,30 @@ async fn get_url_from_runnable( workspace_id: &str, ) -> error::Result { tracing::info!( - "Running {} {} to get websocket URL", + "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).await?; + let args = raw_value_to_args_hashmap(args)?; - if result.starts_with("\"") && result.ends_with("\"") { - Ok(result[1..result.len() - 1].to_string()) - } else { - Err(error::Error::BadConfig(format!( + let result = wait_runnable_result( + path.to_string(), + is_flow, + PushArgsOwned { args, extra: None }, + authed, + db, + workspace_id, + ) + .await?; + + serde_json::from_str::(result.as_str()).map_err(|_| { + error::Error::BadConfig(format!( "{} {} did not return a string", if is_flow { "Flow" } else { "Script" }, - path - ))) - } + path, + )) + }) } impl WebsocketTrigger { @@ -712,11 +724,11 @@ impl WebsocketTrigger { if has_lock.flatten().unwrap_or(false) { tokio::spawn(listen_to_websocket(WebsocketEnum::Trigger(self), db, killpill_rx)); } else { - tracing::info!("Websocket {} already being listened to", self.url); + tracing::info!("WebSocket {} already being listened to", self.url); } }, Err(err) => { - tracing::error!("Error acquiring lock for websocket {}: {:?}", self.path, err); + tracing::error!("Error acquiring lock for WebSocket {}: {:?}", self.path, err); } }; } @@ -737,12 +749,12 @@ impl WebsocketTrigger { self.workspace_id, self.path, ).execute(db).await.ok(); - tracing::info!("Websocket {} changed, disabled, or deleted, stopping...", self.url); + tracing::info!("WebSocket {} changed, disabled, or deleted, stopping...", self.url); return None; } }, Err(err) => { - tracing::warn!("Error updating ping of websocket {}: {:?}", self.url, err); + tracing::warn!("Error updating ping of WebSocket {}: {:?}", self.url, err); } }; @@ -758,11 +770,11 @@ impl WebsocketTrigger { ) .execute(db).await { Ok(_) => { - report_critical_error(format!("Disabling websocket {} because of error: {}", self.url, error), db.clone(), Some(&self.workspace_id), None).await; + report_critical_error(format!("Disabling WebSocket {} because of error: {}", self.url, error), db.clone(), Some(&self.workspace_id), None).await; }, Err(disable_err) => { report_critical_error( - format!("Could not disable websocket {} with err {}, disabling because of error {}", self.path, disable_err, error), + format!("Could not disable WebSocket {} with err {}, disabling because of error {}", self.path, disable_err, error), db.clone(), Some(&self.workspace_id), None, @@ -790,7 +802,7 @@ impl WebsocketTrigger { async fn send_initial_messages( &self, - mut writer: SplitSink>, Message>, + writer: &mut SplitSink>, Message>, db: &DB, ) -> error::Result<()> { let initial_messages: Vec = self @@ -810,7 +822,7 @@ impl WebsocketTrigger { msg }; tracing::info!( - "Sending raw message initial message to websocket {}: {}", + "Sending raw message initial message to WebSocket {}: {}", self.url, msg ); @@ -822,16 +834,18 @@ impl WebsocketTrigger { } InitialMessage::RunnableResult { path, is_flow, args } => { tracing::info!( - "Running {} {} for initial message to websocket {}", + "Running {} {} for initial message to WebSocket {}", if is_flow { "flow" } else { "script" }, path, self.url, ); + let args = raw_value_to_args_hashmap(Some(&args))?; + let result = wait_runnable_result( path.clone(), is_flow, - Some(&args), + PushArgsOwned { args, extra: None }, self.fetch_authed(db).await?, db, &self.workspace_id, @@ -839,17 +853,15 @@ impl WebsocketTrigger { .await?; tracing::info!( - "Sending {} {} result to websocket {}", + "Sending {} {} result to WebSocket {}", if is_flow { "flow" } else { "script" }, path, self.url ); - let result = if result.starts_with("\"") && result.ends_with("\"") { - result[1..result.len() - 1].to_string() - } else { - result - }; + // if the `result` was just a single string, the below removes the surrounding quotes by parsing it as a string. + // it falls back to the original serialized JSON if it doesn't work. + let result = serde_json::from_str::(result.as_str()).unwrap_or(result); writer .send(tokio_tungstenite::tungstenite::Message::Text(result)) @@ -869,11 +881,16 @@ impl WebsocketTrigger { Ok(()) } - async fn handle(&self, db: &DB, args: PushArgsOwned) -> () { - if let Err(err) = run_job(db, self, args).await { + async fn handle( + &self, + db: &DB, + args: PushArgsOwned, + return_message_channels: Option, + ) -> () { + if let Err(err) = run_job(db, self, args, return_message_channels).await { report_critical_error( format!( - "Failed to trigger job from websocket {}: {:?}", + "Failed to trigger job from WebSocket {}: {:?}", self.url, err ), db.clone(), @@ -923,11 +940,11 @@ impl CaptureConfigForWebsocket { if has_lock.flatten().unwrap_or(false) { tokio::spawn(listen_to_websocket(WebsocketEnum::Capture(self), db, killpill_rx)); } else { - tracing::info!("Websocket {} already being listened to", self.trigger_config.url); + tracing::info!("WebSocket {} already being listened to", self.trigger_config.url); } }, Err(err) => { - tracing::error!("Error acquiring lock for capture websocket {}: {:?}", self.path, err); + tracing::error!("Error acquiring lock for capture WebSocket {}: {:?}", self.path, err); } }; } @@ -950,12 +967,12 @@ impl CaptureConfigForWebsocket { self.path, self.is_flow, ).execute(db).await.ok(); - tracing::info!("Websocket capture {} changed, disabled, or deleted, stopping...", self.trigger_config.url); + tracing::info!("WebSocket capture {} changed, disabled, or deleted, stopping...", self.trigger_config.url); return None; } }, Err(err) => { - tracing::warn!("Error updating ping of capture websocket {}: {:?}", self.trigger_config.url, err); + tracing::warn!("Error updating ping of capture WebSocket {}: {:?}", self.trigger_config.url, err); } }; @@ -1021,7 +1038,7 @@ impl CaptureConfigForWebsocket { self.is_flow, ) .execute(db).await { - tracing::error!("Could not disable websocket capture {} ({}) with err {}, disabling because of error {}", self.path, self.workspace_id, err, error); + tracing::error!("Could not disable WebSocket capture {} ({}) with err {}, disabling because of error {}", self.path, self.workspace_id, err, error); } } @@ -1069,6 +1086,20 @@ impl WebsocketEnum { } } +struct ReturnMessageChannels { + send_message_tx: tokio::sync::mpsc::Sender, + killpill_rx: tokio::sync::broadcast::Receiver<()>, +} + +impl Clone for ReturnMessageChannels { + fn clone(&self) -> Self { + Self { + send_message_tx: self.send_message_tx.clone(), + killpill_rx: self.killpill_rx.resubscribe(), + } + } +} + async fn listen_to_websocket( ws: WebsocketEnum, db: DB, @@ -1097,7 +1128,7 @@ async fn listen_to_websocket( return; }, _ = loop_ping(&db, &ws, Some( - "Waiting on runnable to return websocket URL..." + "Waiting on runnable to return WebSocket URL..." )) => { return; }, @@ -1106,7 +1137,7 @@ async fn listen_to_websocket( Ok(url) => Cow::Owned(url), Err(err) => { ws.disable_with_error(&db, format!( - "Error getting websocket URL from runnable after 5 tries: {:?}", + "Error getting WebSocket URL from runnable after 5 tries: {:?}", err ), ) @@ -1116,7 +1147,7 @@ async fn listen_to_websocket( }, } } else { - ws.disable_with_error(&db, format!("Invalid websocket runnable path: {}", url)) + ws.disable_with_error(&db, format!("Invalid WebSocket runnable path: {}", url)) .await; return; } @@ -1135,8 +1166,8 @@ async fn listen_to_websocket( 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(); + tracing::info!("Connected to WebSocket {}", url); + let (mut writer, mut reader) = ws_stream.split(); // send initial messages match &ws { @@ -1149,110 +1180,134 @@ async fn listen_to_websocket( _ = loop_ping(&db, &ws, Some("Sending initial messages...")) => { return; }, - result = ws_trigger.send_initial_messages(writer, &db) => { + result = ws_trigger.send_initial_messages(&mut 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); + tracing::debug!("Initial messages sent successfully to WebSocket {}", url); } } } }, - _ => {} + _ => { + } } - 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 { - tokio_tungstenite::tungstenite::Message::Text(text) => { - tracing::debug!("Received text message from websocket {}: {}", url, text); - let mut should_handle = true; - for filter in &filters { - match filter { - Filter::JsonFilter(JsonFilter { key, value }) => { - let mut deserializer = serde_json::Deserializer::from_str(text.as_str()); - should_handle = match is_value_superset(&mut deserializer, key, &value) { - Ok(filter_match) => { - filter_match - }, - Err(err) => { - tracing::warn!("Error deserializing filter for websocket {}: {:?}", url, err); - false - } - }; - } - } - if !should_handle { - break; - } - } - if should_handle { - - let args = HashMap::from([("msg".to_string(), to_raw_value(&text))]); - let extra = Some(HashMap::from([( - "wm_trigger".to_string(), - to_raw_value(&serde_json::json!({"kind": "websocket", "websocket": { "url": url }})), - )])); - - let args = PushArgsOwned { args, extra }; - match &ws { - WebsocketEnum::Trigger(ws_trigger) => { - ws_trigger.handle(&db, args).await; - }, - WebsocketEnum::Capture(capture) => { - capture.handle(&db, args).await; - }, - } - } - }, - a @ _ => { - tracing::debug!("Received non text-message from websocket {}: {:?}", url, a); - } - } - }, - Err(err) => { - tracing::error!("Error reading from websocket {}: {:?}", url, err); - } - } - } else { - tracing::error!("Websocket {} closed", url); - if let None = ws.update_ping(&db, Some("Websocket closed")).await { - return; - } - return; + let (return_message_channels, message_sender_handle) = match &ws { + WebsocketEnum::Trigger(ws_trigger) if ws_trigger.can_return_message => { + let (send_message_tx, mut rx) = tokio::sync::mpsc::channel::(100); + let w_id = ws_trigger.workspace_id.clone(); + let url = ws_trigger.url.clone(); + let db = db.clone(); + let handle = tokio::spawn(async move { + while let Some(message) = rx.recv().await { + if let Err(err) = writer.send(tokio_tungstenite::tungstenite::Message::Text(message)).await { + report_critical_error(format!("Could not send runnable result to WebSocket {} because of error: {}", url, err), db.clone(), Some(&w_id), None).await; } } - } => { - return; + }); + + let killpill_rx = killpill_rx.resubscribe(); + + let return_message_channels = ReturnMessageChannels { + send_message_tx, + killpill_rx + }; + + (Some(return_message_channels), Some(handle)) + }, + _ => (None, None) + }; + + tokio::select! { + biased; + _ = killpill_rx.recv() => {}, + _ = loop_ping(&db, &ws, None) => {}, + _ = async { + loop { + if let Some(msg) = reader.next().await { + match msg { + Ok(msg) => { + match msg { + tokio_tungstenite::tungstenite::Message::Text(text) => { + tracing::debug!("Received text message from WebSocket {}: {}", url, text); + let mut should_handle = true; + for filter in &filters { + match filter { + Filter::JsonFilter(JsonFilter { key, value }) => { + let mut deserializer = serde_json::Deserializer::from_str(text.as_str()); + should_handle = match is_value_superset(&mut deserializer, key, &value) { + Ok(filter_match) => { + filter_match + }, + Err(err) => { + tracing::warn!("Error deserializing filter for WebSocket {}: {:?}", url, err); + false + } + }; + } + } + if !should_handle { + break; + } + } + if should_handle { + + let args = HashMap::from([("msg".to_string(), to_raw_value(&text))]); + let extra = Some(HashMap::from([( + "wm_trigger".to_string(), + to_raw_value(&serde_json::json!({"kind": "websocket", "websocket": { "url": url }})), + )])); + + let args = PushArgsOwned { args, extra }; + match &ws { + WebsocketEnum::Trigger(ws_trigger) => { + ws_trigger.handle(&db, args, return_message_channels.clone()).await; + }, + WebsocketEnum::Capture(capture) => { + capture.handle(&db, args).await; + }, + } + } + }, + a @ _ => { + tracing::debug!("Received non text-message from WebSocket {}: {:?}", url, a); + } + } + }, + Err(err) => { + tracing::error!("Error reading from WebSocket {}: {:?}", url, err); + } + } + } else { + tracing::error!("WebSocket {} closed", url); + ws.update_ping(&db, Some("WebSocket closed")).await; + break; + } } - } + } => {} + } + // make sure to stop return message handler + if let Some(message_sender_handle) = message_sender_handle { + message_sender_handle.abort(); } } 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; - } + tracing::error!("Error connecting to WebSocket {}: {:?}", url, err); + ws.update_ping(&db, Some(err.to_string().as_str())).await; } } } } } -async fn run_job(db: &DB, trigger: &WebsocketTrigger, args: PushArgsOwned) -> anyhow::Result<()> { +async fn run_job( + db: &DB, + trigger: &WebsocketTrigger, + args: PushArgsOwned, + return_message_channels: Option, +) -> anyhow::Result<()> { let authed = fetch_api_authed( trigger.edited_by.clone(), trigger.email.clone(), @@ -1262,34 +1317,73 @@ async fn run_job(db: &DB, trigger: &WebsocketTrigger, args: PushArgsOwned) -> an ) .await?; - let user_db = UserDB::new(db.clone()); + if let Some(ReturnMessageChannels { send_message_tx, mut killpill_rx }) = + return_message_channels + { + let db_ = db.clone(); + let url = trigger.url.clone(); + let script_path = trigger.script_path.clone(); + let is_flow = trigger.is_flow; + let w_id = trigger.workspace_id.clone(); + let handle_response_f = async move { + tokio::select! { + _ = killpill_rx.recv() => { + return; + }, + result = wait_runnable_result( + script_path, + is_flow, + args, + authed, + &db_, + &w_id, + ) => { + if let Ok(result) = result { + // only send the result if it's not null + if result != "null" { + tracing::info!("Sending job result to WebSocket {}", url); + // if the `result` was just a single string, the below removes the surrounding quotes by parsing it as a string. + // it falls back to the original serialized JSON if it doesn't work. + let result = serde_json::from_str::(result.as_str()).unwrap_or(result); + if let Err(err) = send_message_tx.send(result).await { + report_critical_error(format!("Could not send runnable result to WebSocket {} because of error: {}", url, err), db_.clone(), Some(&w_id), None).await; + } + } + } + } + }; + }; - let run_query = RunJobQuery::default(); - - if trigger.is_flow { - run_flow_by_path_inner( - authed, - db.clone(), - user_db, - trigger.workspace_id.clone(), - StripPath(trigger.script_path.to_owned()), - run_query, - args, - None, - ) - .await?; + tokio::spawn(handle_response_f); } else { - run_script_by_path_inner( - authed, - db.clone(), - user_db, - trigger.workspace_id.clone(), - StripPath(trigger.script_path.to_owned()), - run_query, - args, - None, - ) - .await?; + let user_db = UserDB::new(db.clone()); + let run_query = RunJobQuery::default(); + let runnable_path = StripPath(trigger.script_path.to_owned()); + if trigger.is_flow { + run_flow_by_path_inner( + authed, + db.clone(), + user_db, + trigger.workspace_id.clone(), + runnable_path, + run_query, + args, + None, + ) + .await?; + } else { + run_script_by_path_inner( + authed, + db.clone(), + user_db, + trigger.workspace_id.clone(), + runnable_path, + run_query, + args, + None, + ) + .await?; + } } Ok(()) diff --git a/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte b/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte index c5d4938828..e08f338c6e 100644 --- a/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte +++ b/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte @@ -64,7 +64,7 @@ - Websockets + WebSockets diff --git a/frontend/src/lib/components/sidebar/SidebarContent.svelte b/frontend/src/lib/components/sidebar/SidebarContent.svelte index 7edc576bc8..cf82e1aaee 100644 --- a/frontend/src/lib/components/sidebar/SidebarContent.svelte +++ b/frontend/src/lib/components/sidebar/SidebarContent.svelte @@ -95,7 +95,7 @@ kind: 'http' }, { - label: 'Websockets', + label: 'WebSockets', href: '/websocket_triggers', icon: Unplug, disabled: $userStore?.operator, diff --git a/frontend/src/lib/components/triggers/TestTriggerConnection.svelte b/frontend/src/lib/components/triggers/TestTriggerConnection.svelte index 8c56f3627b..511b6c3f9d 100644 --- a/frontend/src/lib/components/triggers/TestTriggerConnection.svelte +++ b/frontend/src/lib/components/triggers/TestTriggerConnection.svelte @@ -13,7 +13,7 @@ export let args: Record const kindToName: { [key: string]: string } = { - websocket: 'Websocket', + websocket: 'WebSocket', nats: 'NATS server(s)', kafka: 'Kafka broker(s)' } @@ -47,7 +47,6 @@ await promise sendUserToast(`Successfully connected to ${kindToName[kind]}`) } catch (err) { - if (!promise?.isCancelled) { sendUserToast(`Error testing ${kindToName[kind]}: ${err?.body ?? 'Unknown error'}`, true) } diff --git a/frontend/src/lib/components/triggers/TriggersEditor.svelte b/frontend/src/lib/components/triggers/TriggersEditor.svelte index 922ac7f066..29d739036b 100644 --- a/frontend/src/lib/components/triggers/TriggersEditor.svelte +++ b/frontend/src/lib/components/triggers/TriggersEditor.svelte @@ -53,7 +53,7 @@ Webhooks Schedules HTTP - Websockets + WebSockets Postgres = { http: 'New custom HTTP route', - websocket: 'New websocket trigger', + websocket: 'New WebSocket trigger', webhook: 'Webhook', - kafka: 'New kafka trigger', + kafka: 'New Kafka trigger', email: 'Email trigger', nats: 'NATS trigger' } diff --git a/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte index 00883a8ff0..8aeedb77e7 100644 --- a/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte @@ -46,7 +46,7 @@ dirtyPath = false await loadTrigger() } catch (err) { - sendUserToast(`Could not load kafka trigger: ${err}`, true) + sendUserToast(`Could not load Kafka trigger: ${err}`, true) } finally { drawerLoading = false } @@ -154,9 +154,9 @@ @@ -173,7 +173,7 @@ workspace: $workspaceStore ?? '', requestBody: { enabled: e.detail } }) - sendUserToast(`${e.detail ? 'enabled' : 'disabled'} kafka trigger ${initialPath}`) + sendUserToast(`${e.detail ? 'enabled' : 'disabled'} Kafka trigger ${initialPath}`) }} /> diff --git a/frontend/src/lib/components/triggers/kafka/KafkaTriggersPanel.svelte b/frontend/src/lib/components/triggers/kafka/KafkaTriggersPanel.svelte index 08ac6c9b64..5e710676d1 100644 --- a/frontend/src/lib/components/triggers/kafka/KafkaTriggersPanel.svelte +++ b/frontend/src/lib/components/triggers/kafka/KafkaTriggersPanel.svelte @@ -51,7 +51,7 @@ }) $triggersCount = { ...($triggersCount ?? {}), kafka_count: kafkaTriggers?.length } } catch (e) { - console.error('impossible to load kafka triggers', e) + console.error('impossible to load Kafka triggers', e) } } @@ -108,7 +108,7 @@ {#if kafkaTriggers}
{#if kafkaTriggers.length == 0} -
No kafka triggers
+
No Kafka triggers
{:else}
{#each kafkaTriggers as kafkaTrigger (kafkaTrigger.path)} diff --git a/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte b/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte index aac583c09d..649b7bcff1 100644 --- a/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte +++ b/frontend/src/lib/components/triggers/websocket/WebsocketEditorConfigSection.svelte @@ -94,7 +94,7 @@ bind:captureTable /> {/if} -
+
| undefined = {} + let can_return_message = false let dirtyPath = false let can_write = true let drawerLoading = true @@ -92,6 +93,7 @@ initial_messages = [] url_runnable_args = defaultValues?.url_runnable_args ?? {} dirtyPath = false + can_return_message = false } finally { drawerLoading = false } @@ -112,6 +114,7 @@ filters = s.filters initial_messages = s.initial_messages ?? [] url_runnable_args = s.url_runnable_args + can_return_message = s.can_return_message can_write = canWrite(s.path, s.extra_perms, $userStore) } @@ -168,7 +171,8 @@ url, filters, initial_messages, - url_runnable_args + url_runnable_args, + can_return_message } }) sendUserToast(`Websocket trigger ${path} updated`) @@ -183,7 +187,8 @@ enabled: true, filters, initial_messages, - url_runnable_args + url_runnable_args, + can_return_message } }) sendUserToast(`Websocket trigger ${path} created`) @@ -202,9 +207,9 @@ @@ -250,7 +255,7 @@ {#if edit} Changes can take up to 30 seconds to take effect. {:else} - New websocket triggers can take up to 30 seconds to start listening. + New WebSocket triggers can take up to 30 seconds to start listening. {/if}
@@ -277,22 +282,36 @@ bind:isValid /> -
-

- Pick a script or flow to be triggered -

-
- +
+
+

+ Pick a script or flow to be triggered +

+
+ +
+ + { + can_return_message = !can_return_message + }} + options={{ + right: 'Send runnable result', + rightTooltip: + 'Whether the runnable result should be sent as a message to the websocket server when not null.' + }} + />
diff --git a/frontend/src/lib/components/triggers/websocket/WebsocketTriggersPanel.svelte b/frontend/src/lib/components/triggers/websocket/WebsocketTriggersPanel.svelte index b53862a82c..7b633a1262 100644 --- a/frontend/src/lib/components/triggers/websocket/WebsocketTriggersPanel.svelte +++ b/frontend/src/lib/components/triggers/websocket/WebsocketTriggersPanel.svelte @@ -64,13 +64,13 @@ {#if isCloudHosted()} - Websocket triggers are disabled in the multi-tenant cloud. + WebSocket triggers are disabled in the multi-tenant cloud. {:else}
- Websocket triggers allow real-time bidirectional communication between your scripts/flows and - external systems. Each trigger creates a unique websocket endpoint. + WebSocket triggers allow real-time bidirectional communication between your scripts/flows and + external systems. Each trigger creates a unique WebSocket endpoint. {#if !newItem} -
+
{#if wsTriggers} {#if wsTriggers.length == 0} -
No WS triggers
+
No WebSocket triggers
{:else}
{#each wsTriggers as wsTriggers (wsTriggers.path)} diff --git a/frontend/src/routes/(root)/(logged)/kafka_triggers/+page.svelte b/frontend/src/routes/(root)/(logged)/kafka_triggers/+page.svelte index 13c9cfd1d9..d4dba36df7 100644 --- a/frontend/src/routes/(root)/(logged)/kafka_triggers/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/kafka_triggers/+page.svelte @@ -81,7 +81,7 @@ }) } catch (err) { sendUserToast( - `Cannot ` + (enabled ? 'enable' : 'disable') + ` kafka trigger: ${err.body}`, + `Cannot ` + (enabled ? 'enable' : 'disable') + ` Kafka trigger: ${err.body}`, true ) } finally { @@ -208,223 +208,230 @@ f={(x) => (x.summary ?? '') + ' ' + x.path + ' (' + x.script_path + ')'} /> -{#if $userStore?.operator && $workspaceStore && !$userWorkspaces.find(_ => _.id === $workspaceStore)?.operator_settings?.triggers} - +{#if $userStore?.operator && $workspaceStore && !$userWorkspaces.find((_) => _.id === $workspaceStore)?.operator_settings?.triggers} + {:else} - - - - + + + + - {#if isCloudHosted()} - - Kafka triggers are disabled in the multi-tenant cloud. - -
- {/if} -
-
- -
-
Filter by path of
- - - - + {#if isCloudHosted()} + + Kafka triggers are disabled in the multi-tenant cloud. + +
+ {/if} +
+
+ +
+
Filter by path of
+ + + + +
+ + +
+ {#if $userStore?.is_super_admin && $userStore.username.includes('@')} + + {:else if $userStore?.is_admin || $userStore?.is_super_admin} + + {/if} +
- + {#if loading} + {#each new Array(6) as _} + + {/each} + {:else if !triggers?.length} +
No Kafka triggers
+ {:else if items?.length} +
+ {#each items.slice(0, nbDisplayed) as { path, edited_by, edited_at, script_path, is_flow, kafka_resource_path, topics, extra_perms, canWrite, marked, server_id, error, last_server_ping, enabled } (path)} + {@const href = `${is_flow ? '/flows/get' : '/scripts/get'}/${script_path}`} + {@const ping = last_server_ping ? new Date(last_server_ping) : undefined} + {@const pinging = ping && ping.getTime() > new Date().getTime() - 15 * 1000} -
- {#if $userStore?.is_super_admin && $userStore.username.includes('@')} - - {:else if $userStore?.is_admin || $userStore?.is_super_admin} - - {/if} -
-
- {#if loading} - {#each new Array(6) as _} - - {/each} - {:else if !triggers?.length} -
No kafka triggers
- {:else if items?.length} -
- {#each items.slice(0, nbDisplayed) as { path, edited_by, edited_at, script_path, is_flow, kafka_resource_path, topics, extra_perms, canWrite, marked, server_id, error, last_server_ping, enabled } (path)} - {@const href = `${is_flow ? '/flows/get' : '/scripts/get'}/${script_path}`} - {@const ping = last_server_ping ? new Date(last_server_ping) : undefined} - {@const pinging = ping && ping.getTime() > new Date().getTime() - 15 * 1000} - -
-
- + > +
+ - kafkaTriggerEditor?.openEdit(path, is_flow)} - class="min-w-0 grow hover:underline decoration-gray-400" - > -
- {#if marked} - - {@html marked} - - {:else} - {kafka_resource_path} - {topics.join(', ')} + kafkaTriggerEditor?.openEdit(path, is_flow)} + class="min-w-0 grow hover:underline decoration-gray-400" + > +
+ {#if marked} + + {@html marked} + + {:else} + {kafka_resource_path} - {topics.join(', ')} + {/if} +
+
+ {path} +
+
+ runnable: {script_path} +
+
+ + + +
+ {#if (enabled && (!pinging || error)) || (!enabled && error) || (enabled && !server_id)} + + + + + +
+ {#if enabled} + {#if !server_id} + Consumer is starting... + {:else} + Consumer is not connected{error ? ': ' + error : ''} + {/if} + {:else} + Consumer was disabled because of an error: {error} + {/if} +
+
+ {:else if enabled} + + + + +
Consumer is connected
+
{/if}
-
- {path} -
-
- runnable: {script_path} -
- - - -
- {#if (enabled && (!pinging || error)) || (!enabled && error) || (enabled && !server_id)} - - - - - -
- {#if enabled} - {#if !server_id} - Consumer is starting... - {:else} - Consumer is not connected{error ? ': ' + error : ''} - {/if} - {:else} - Consumer was disabled because of an error: {error} - {/if} -
-
- {:else if enabled} - - - - -
Consumer is connected
-
- {/if} -
- - { - setTriggerEnabled(path, e.detail) - }} - /> - -
- - { - goto(href) - } - }, - { - displayName: 'Delete', - type: 'delete', - icon: Trash, - disabled: !canWrite, - action: async () => { - await KafkaTriggerService.deleteKafkaTrigger({ - workspace: $workspaceStore ?? '', - path - }) - loadTriggers() - } - }, - { - displayName: canWrite ? 'Edit' : 'View', - icon: canWrite ? Pen : Eye, - action: () => { - kafkaTriggerEditor?.openEdit(path, is_flow) - } - }, - { - displayName: 'Audit logs', - icon: Eye, - href: `${base}/audit_logs?resource=${path}` - }, - { - displayName: canWrite ? 'Share' : 'See Permissions', - icon: Share, - action: () => { - shareModal.openDrawer(path, 'kafka_trigger') - } - } - ]} + { + setTriggerEnabled(path, e.detail) + }} /> + +
+ + { + goto(href) + } + }, + { + displayName: 'Delete', + type: 'delete', + icon: Trash, + disabled: !canWrite, + action: async () => { + await KafkaTriggerService.deleteKafkaTrigger({ + workspace: $workspaceStore ?? '', + path + }) + loadTriggers() + } + }, + { + displayName: canWrite ? 'Edit' : 'View', + icon: canWrite ? Pen : Eye, + action: () => { + kafkaTriggerEditor?.openEdit(path, is_flow) + } + }, + { + displayName: 'Audit logs', + icon: Eye, + href: `${base}/audit_logs?resource=${path}` + }, + { + displayName: canWrite ? 'Share' : 'See Permissions', + icon: Share, + action: () => { + shareModal.openDrawer(path, 'kafka_trigger') + } + } + ]} + /> +
-
-
-
edited by {edited_by}
the {displayDate(edited_at)}
+
edited by {edited_by}
the {displayDate(edited_at)}
-
- {/each} -
- {:else} - + > +
+ {/each} +
+ {:else} + + {/if} +
+ {#if items && items?.length > 15 && nbDisplayed < items.length} + {nbDisplayed} items out of {items.length} + {/if} -
- {#if items && items?.length > 15 && nbDisplayed < items.length} - {nbDisplayed} items out of {items.length} - - {/if} - + {/if} {#if isCloudHosted()} - Websocket triggers are disabled in the multi-tenant cloud. + WebSocket triggers are disabled in the multi-tenant cloud.
{/if} @@ -314,12 +314,12 @@
{#if enabled} {#if !server_id} - Websocket is starting... + WebSocket is starting... {:else} - Websocket is not connected{error ? ': ' + error : ''} + WebSocket is not connected{error ? ': ' + error : ''} {/if} {:else} - Websocket was disabled because of an error: {error} + WebSocket was disabled because of an error: {error} {/if}
@@ -329,7 +329,7 @@
- Websocket is connected{!server_id ? ' (shutting down...)' : ''} + WebSocket is connected{!server_id ? ' (shutting down...)' : ''}
{/if}