feat: websocket trigger allow returning messages (#5168)

* feat: websocket trigger allow returning messages

* handle early return
This commit is contained in:
HugoCasa
2025-01-29 21:17:36 +01:00
committed by GitHub
parent cfcd84be3f
commit 487c273bfb
21 changed files with 740 additions and 614 deletions
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -0,0 +1 @@
ALTER TABLE websocket_trigger DROP COLUMN can_return_message;
@@ -0,0 +1 @@
ALTER TABLE websocket_trigger ADD COLUMN can_return_message BOOLEAN NOT NULL DEFAULT FALSE;
+12 -1
View File
@@ -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
+103 -93
View File
@@ -3527,13 +3527,14 @@ pub struct WindmillCompositeResult {
windmill_headers: Option<HashMap<String, String>>,
result: Option<Box<RawValue>>,
}
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<String>,
username: &str,
) -> error::Result<Response> {
) -> error::Result<(Box<RawValue>, 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::<WindmillCompositeResult>(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::<String>(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<String>,
username: &str,
) -> error::Result<Response> {
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::<WindmillCompositeResult>(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::<String>(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
+318 -224
View File
@@ -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<Box<RawValue>>,
initial_messages: Option<Vec<Box<RawValue>>>,
url_runnable_args: Option<Box<RawValue>>,
can_return_message: bool,
}
#[derive(Deserialize)]
@@ -101,6 +104,7 @@ pub struct WebsocketTrigger {
filters: Vec<SqlxJson<Box<RawValue>>>,
initial_messages: Option<Vec<SqlxJson<Box<RawValue>>>>,
url_runnable_args: Option<SqlxJson<Box<RawValue>>>,
can_return_message: bool,
}
#[derive(Deserialize)]
@@ -112,6 +116,7 @@ struct EditWebsocketTrigger {
filters: Vec<Box<RawValue>>,
initial_messages: Option<Vec<Box<RawValue>>>,
url_runnable_args: Option<Box<RawValue>>,
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<SqlxJson<Box<RawValue>>>,
&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<RawValue>>,
authed: ApiAuthed,
db: &DB,
workspace_id: &str,
) -> error::Result<String> {
let user_db = UserDB::new(db.clone());
) -> error::Result<HashMap<String, Box<RawValue>>> {
let args = if let Some(args) = args {
serde_json::from_str::<Option<HashMap<String, Box<RawValue>>>>(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<String> {
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<SqlxJson<Box<RawValue>>>,
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<String> {
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::<String>(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<WebSocketStream<MaybeTlsStream<TcpStream>>, Message>,
writer: &mut SplitSink<WebSocketStream<MaybeTlsStream<TcpStream>>, Message>,
db: &DB,
) -> error::Result<()> {
let initial_messages: Vec<InitialMessage> = 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::<String>(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<ReturnMessageChannels>,
) -> () {
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<String>,
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::<String>(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<ReturnMessageChannels>,
) -> 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::<String>(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(())
@@ -64,7 +64,7 @@
<Tab value="websockets">
<span class="flex flex-row gap-2 items-center text-xs">
<Unplug size={12} />
Websockets
WebSockets
</span>
</Tab>
<Tab value="postgres">
@@ -95,7 +95,7 @@
kind: 'http'
},
{
label: 'Websockets',
label: 'WebSockets',
href: '/websocket_triggers',
icon: Unplug,
disabled: $userStore?.operator,
@@ -13,7 +13,7 @@
export let args: Record<string, any>
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)
}
@@ -53,7 +53,7 @@
<Tab value="webhooks" selectedClass="text-primary font-semibold">Webhooks</Tab>
<Tab value="schedules" selectedClass="text-primary text-sm font-semibold">Schedules</Tab>
<Tab value="routes" selectedClass="text-primary text-sm font-semibold">HTTP</Tab>
<Tab value="websockets" selectedClass="text-primary text-sm font-semibold">Websockets</Tab>
<Tab value="websockets" selectedClass="text-primary text-sm font-semibold">WebSockets</Tab>
<Tab value="postgres" selectedClass="text-primary text-sm font-semibold">Postgres</Tab>
<Tab
value="kafka"
@@ -24,9 +24,9 @@
const captureTypeLabels: Record<CaptureTriggerKind, string> = {
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'
}
@@ -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 @@
<DrawerContent
title={edit
? can_write
? `Edit kafka trigger ${initialPath}`
? `Edit Kafka trigger ${initialPath}`
: `Kafka trigger ${initialPath}`
: 'New kafka trigger'}
: 'New Kafka trigger'}
on:close={drawer.closeDrawer}
>
<svelte:fragment slot="actions">
@@ -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}`)
}}
/>
</div>
@@ -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}
<Section label="Kafka triggers">
{#if kafkaTriggers.length == 0}
<div class="text-xs text-secondary text-center"> No kafka triggers </div>
<div class="text-xs text-secondary text-center"> No Kafka triggers </div>
{:else}
<div class="flex flex-col divide-y pt-2">
{#each kafkaTriggers as kafkaTrigger (kafkaTrigger.path)}
@@ -94,7 +94,7 @@
bind:captureTable
/>
{/if}
<Section label="Websocket" {headless}>
<Section label="WebSocket" {headless}>
<div class="mb-2">
<ToggleButtonGroup
selected={url?.startsWith('$') ? 'runnable' : 'static'}
@@ -45,6 +45,7 @@
}[] = []
let initial_messages: WebsocketTriggerInitialMessage[] = []
let url_runnable_args: Record<string, any> | 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 @@
<DrawerContent
title={edit
? can_write
? `Edit WS trigger ${initialPath}`
: `WS trigger ${initialPath}`
: 'New WS trigger'}
? `Edit WebSocket trigger ${initialPath}`
: `WebSocket trigger ${initialPath}`
: 'New WebSocket trigger'}
on:close={drawer.closeDrawer}
>
<svelte:fragment slot="actions">
@@ -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}
</Alert>
<div class="flex flex-col gap-12 mt-6">
@@ -277,22 +282,36 @@
bind:isValid
/>
<Section label="Runnable">
<p class="text-xs mb-1 text-tertiary">
Pick a script or flow to be triggered<Required required={true} />
</p>
<div class="flex flex-row mb-2">
<ScriptPicker
disabled={fixedScriptPath != '' || !can_write}
initialPath={fixedScriptPath || initialScriptPath}
kinds={['script']}
allowFlow={true}
bind:itemKind
bind:scriptPath={script_path}
allowRefresh={can_write}
allowEdit={!$userStore?.operator}
/>
<Section label="Runnable" class="flex flex-col gap-4">
<div>
<p class="text-xs mb-1 text-tertiary">
Pick a script or flow to be triggered<Required required={true} />
</p>
<div class="flex flex-row mb-2">
<ScriptPicker
disabled={fixedScriptPath != '' || !can_write}
initialPath={fixedScriptPath || initialScriptPath}
kinds={['script']}
allowFlow={true}
bind:itemKind
bind:scriptPath={script_path}
allowRefresh={can_write}
allowEdit={!$userStore?.operator}
/>
</div>
</div>
<Toggle
checked={can_return_message}
on:change={() => {
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.'
}}
/>
</Section>
<Section label="Initial messages">
@@ -64,13 +64,13 @@
{#if isCloudHosted()}
<Alert title="Not compatible with multi-tenant cloud" type="warning" size="xs">
Websocket triggers are disabled in the multi-tenant cloud.
WebSocket triggers are disabled in the multi-tenant cloud.
</Alert>
{:else}
<div class="flex flex-col gap-4">
<Description link="https://www.windmill.dev/docs/core_concepts/websocket_triggers">
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.
</Description>
<TriggersEditorSection
on:applyArgs
@@ -91,11 +91,11 @@
/>
{#if !newItem}
<Section label="Websockets">
<Section label="WebSockets">
<div class="flex flex-col gap-4">
{#if wsTriggers}
{#if wsTriggers.length == 0}
<div class="text-xs text-secondary text-center"> No WS triggers </div>
<div class="text-xs text-secondary text-center"> No WebSocket triggers </div>
{:else}
<div class="flex flex-col divide-y pt-2">
{#each wsTriggers as wsTriggers (wsTriggers.path)}
@@ -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}
<div class="bg-red-100 border-l-4 border-red-600 text-orange-700 p-4 m-4 mt-12" role="alert">
<p class="font-bold">Unauthorized</p>
<p>Page not available for operators</p>
</div>
{#if $userStore?.operator && $workspaceStore && !$userWorkspaces.find((_) => _.id === $workspaceStore)?.operator_settings?.triggers}
<div class="bg-red-100 border-l-4 border-red-600 text-orange-700 p-4 m-4 mt-12" role="alert">
<p class="font-bold">Unauthorized</p>
<p>Page not available for operators</p>
</div>
{:else}
<CenteredPage>
<PageHeader
title="Kafka triggers"
tooltip="Windmill can consume kafka events and trigger scripts or flows based on them."
>
<Button size="md" startIcon={{ icon: Plus }} on:click={() => kafkaTriggerEditor.openNew(false)}>
New&nbsp;kafka trigger
</Button>
</PageHeader>
<CenteredPage>
<PageHeader
title="Kafka triggers"
tooltip="Windmill can consume kafka events and trigger scripts or flows based on them."
>
<Button
size="md"
startIcon={{ icon: Plus }}
on:click={() => kafkaTriggerEditor.openNew(false)}
>
New&nbsp;Kafka trigger
</Button>
</PageHeader>
{#if isCloudHosted()}
<Alert title="Not compatible with multi-tenant cloud" type="warning">
Kafka triggers are disabled in the multi-tenant cloud.
</Alert>
<div class="py-4" />
{/if}
<div class="w-full h-full flex flex-col">
<div class="w-full pb-4 pt-6">
<input
type="text"
placeholder="Search kafka triggers"
bind:value={filter}
class="search-item"
/>
<div class="flex flex-row items-center gap-2 mt-6">
<div class="text-sm shrink-0"> Filter by path of </div>
<ToggleButtonGroup bind:selected={selectedFilterKind}>
<ToggleButton small value="trigger" label="Kafka trigger" icon={KafkaIcon} />
<ToggleButton small value="script_flow" label="Script/Flow" icon={Code} />
</ToggleButtonGroup>
{#if isCloudHosted()}
<Alert title="Not compatible with multi-tenant cloud" type="warning">
Kafka triggers are disabled in the multi-tenant cloud.
</Alert>
<div class="py-4" />
{/if}
<div class="w-full h-full flex flex-col">
<div class="w-full pb-4 pt-6">
<input
type="text"
placeholder="Search Kafka triggers"
bind:value={filter}
class="search-item"
/>
<div class="flex flex-row items-center gap-2 mt-6">
<div class="text-sm shrink-0"> Filter by path of </div>
<ToggleButtonGroup bind:selected={selectedFilterKind}>
<ToggleButton small value="trigger" label="Kafka trigger" icon={KafkaIcon} />
<ToggleButton small value="script_flow" label="Script/Flow" icon={Code} />
</ToggleButtonGroup>
</div>
<ListFilters syncQuery bind:selectedFilter={ownerFilter} filters={owners} />
<div class="flex flex-row items-center justify-end gap-4">
{#if $userStore?.is_super_admin && $userStore.username.includes('@')}
<Toggle size="xs" bind:checked={filterUserFolders} options={{ right: 'Only f/*' }} />
{:else if $userStore?.is_admin || $userStore?.is_super_admin}
<Toggle
size="xs"
bind:checked={filterUserFolders}
options={{ right: `Only u/${$userStore.username} and f/*` }}
/>
{/if}
</div>
</div>
<ListFilters syncQuery bind:selectedFilter={ownerFilter} filters={owners} />
{#if loading}
{#each new Array(6) as _}
<Skeleton layout={[[6], 0.4]} />
{/each}
{:else if !triggers?.length}
<div class="text-center text-sm text-tertiary mt-2"> No Kafka triggers </div>
{:else if items?.length}
<div class="border rounded-md divide-y">
{#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}
<div class="flex flex-row items-center justify-end gap-4">
{#if $userStore?.is_super_admin && $userStore.username.includes('@')}
<Toggle size="xs" bind:checked={filterUserFolders} options={{ right: 'Only f/*' }} />
{:else if $userStore?.is_admin || $userStore?.is_super_admin}
<Toggle
size="xs"
bind:checked={filterUserFolders}
options={{ right: `Only u/${$userStore.username} and f/*` }}
/>
{/if}
</div>
</div>
{#if loading}
{#each new Array(6) as _}
<Skeleton layout={[[6], 0.4]} />
{/each}
{:else if !triggers?.length}
<div class="text-center text-sm text-tertiary mt-2"> No kafka triggers </div>
{:else if items?.length}
<div class="border rounded-md divide-y">
{#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}
<div
class="hover:bg-surface-hover w-full items-center px-4 py-2 gap-4 first-of-type:!border-t-0
<div
class="hover:bg-surface-hover w-full items-center px-4 py-2 gap-4 first-of-type:!border-t-0
first-of-type:rounded-t-md last-of-type:rounded-b-md flex flex-col"
>
<div class="w-full flex gap-5 items-center">
<RowIcon kind={is_flow ? 'flow' : 'script'} />
>
<div class="w-full flex gap-5 items-center">
<RowIcon kind={is_flow ? 'flow' : 'script'} />
<a
href="#{path}"
on:click={() => kafkaTriggerEditor?.openEdit(path, is_flow)}
class="min-w-0 grow hover:underline decoration-gray-400"
>
<div class="text-primary flex-wrap text-left text-md font-semibold mb-1 truncate">
{#if marked}
<span class="text-xs">
{@html marked}
</span>
{:else}
{kafka_resource_path} - {topics.join(', ')}
<a
href="#{path}"
on:click={() => kafkaTriggerEditor?.openEdit(path, is_flow)}
class="min-w-0 grow hover:underline decoration-gray-400"
>
<div class="text-primary flex-wrap text-left text-md font-semibold mb-1 truncate">
{#if marked}
<span class="text-xs">
{@html marked}
</span>
{:else}
{kafka_resource_path} - {topics.join(', ')}
{/if}
</div>
<div class="text-secondary text-xs truncate text-left font-light">
{path}
</div>
<div class="text-secondary text-xs truncate text-left font-light">
runnable: {script_path}
</div>
</a>
<div class="hidden lg:flex flex-row gap-1 items-center">
<SharedBadge {canWrite} extraPerms={extra_perms} />
</div>
<div class="w-10">
{#if (enabled && (!pinging || error)) || (!enabled && error) || (enabled && !server_id)}
<Popover notClickable>
<span class="flex h-4 w-4">
<Circle
class="text-red-600 animate-ping absolute inline-flex fill-current"
size={12}
/>
<Circle class="text-red-600 relative inline-flex fill-current" size={12} />
</span>
<div slot="text">
{#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}
</div>
</Popover>
{:else if enabled}
<Popover notClickable>
<span class="flex h-4 w-4">
<Circle
class="text-green-600 relative inline-flex fill-current"
size={12}
/>
</span>
<div slot="text"> Consumer is connected </div>
</Popover>
{/if}
</div>
<div class="text-secondary text-xs truncate text-left font-light">
{path}
</div>
<div class="text-secondary text-xs truncate text-left font-light">
runnable: {script_path}
</div>
</a>
<div class="hidden lg:flex flex-row gap-1 items-center">
<SharedBadge {canWrite} extraPerms={extra_perms} />
</div>
<div class="w-10">
{#if (enabled && (!pinging || error)) || (!enabled && error) || (enabled && !server_id)}
<Popover notClickable>
<span class="flex h-4 w-4">
<Circle
class="text-red-600 animate-ping absolute inline-flex fill-current"
size={12}
/>
<Circle class="text-red-600 relative inline-flex fill-current" size={12} />
</span>
<div slot="text">
{#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}
</div>
</Popover>
{:else if enabled}
<Popover notClickable>
<span class="flex h-4 w-4">
<Circle class="text-green-600 relative inline-flex fill-current" size={12} />
</span>
<div slot="text"> Consumer is connected </div>
</Popover>
{/if}
</div>
<Toggle
checked={enabled}
disabled={!canWrite}
on:change={(e) => {
setTriggerEnabled(path, e.detail)
}}
/>
<div class="flex gap-2 items-center justify-end">
<Button
on:click={() => kafkaTriggerEditor?.openEdit(path, is_flow)}
size="xs"
startIcon={canWrite
? { icon: Pen }
: {
icon: Eye
}}
color="dark"
>
{canWrite ? 'Edit' : 'View'}
</Button>
<Dropdown
items={[
{
displayName: `View ${is_flow ? 'Flow' : 'Script'}`,
icon: Eye,
action: () => {
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')
}
}
]}
<Toggle
checked={enabled}
disabled={!canWrite}
on:change={(e) => {
setTriggerEnabled(path, e.detail)
}}
/>
<div class="flex gap-2 items-center justify-end">
<Button
on:click={() => kafkaTriggerEditor?.openEdit(path, is_flow)}
size="xs"
startIcon={canWrite
? { icon: Pen }
: {
icon: Eye
}}
color="dark"
>
{canWrite ? 'Edit' : 'View'}
</Button>
<Dropdown
items={[
{
displayName: `View ${is_flow ? 'Flow' : 'Script'}`,
icon: Eye,
action: () => {
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')
}
}
]}
/>
</div>
</div>
</div>
<div class="w-full flex justify-between items-baseline">
<div
class="flex flex-wrap text-[0.7em] text-tertiary gap-1 items-center justify-end truncate pr-2"
><div class="truncate">edited by {edited_by}</div><div class="truncate"
>the {displayDate(edited_at)}</div
<div class="w-full flex justify-between items-baseline">
<div
class="flex flex-wrap text-[0.7em] text-tertiary gap-1 items-center justify-end truncate pr-2"
><div class="truncate">edited by {edited_by}</div><div class="truncate"
>the {displayDate(edited_at)}</div
></div
></div
></div
>
</div>
{/each}
</div>
{:else}
<NoItemFound />
>
</div>
{/each}
</div>
{:else}
<NoItemFound />
{/if}
</div>
{#if items && items?.length > 15 && nbDisplayed < items.length}
<span class="text-xs"
>{nbDisplayed} items out of {items.length}
<button class="ml-4" on:click={() => (nbDisplayed += 30)}>load 30 more</button></span
>
{/if}
</div>
{#if items && items?.length > 15 && nbDisplayed < items.length}
<span class="text-xs"
>{nbDisplayed} items out of {items.length}
<button class="ml-4" on:click={() => (nbDisplayed += 30)}>load 30 more</button></span
>
{/if}
</CenteredPage>
</CenteredPage>
{/if}
<ShareModal
@@ -1,5 +1,5 @@
export function load() {
return {
stuff: { title: 'WS triggers' }
stuff: { title: 'WebSocket triggers' }
}
}
@@ -209,21 +209,21 @@
<CenteredPage>
<PageHeader
title="Websocket triggers"
tooltip="Windmill can listen to websocket events and trigger scripts or flows based on them."
title="WebSocket triggers"
tooltip="Windmill can listen to WebSocket events and trigger scripts or flows based on them."
>
<Button
size="md"
startIcon={{ icon: Plus }}
on:click={() => websocketTriggerEditor.openNew(false)}
>
New&nbsp;WS trigger
New&nbsp;WebSocket trigger
</Button>
</PageHeader>
{#if isCloudHosted()}
<Alert title="Not compatible with multi-tenant cloud" type="warning">
Websocket triggers are disabled in the multi-tenant cloud.
WebSocket triggers are disabled in the multi-tenant cloud.
</Alert>
<div class="py-4" />
{/if}
@@ -314,12 +314,12 @@
<div slot="text">
{#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}
</div>
</Popover>
@@ -329,7 +329,7 @@
<Circle class="text-green-600 relative inline-flex fill-current" size={12} />
</span>
<div slot="text">
Websocket is connected{!server_id ? ' (shutting down...)' : ''}
WebSocket is connected{!server_id ? ' (shutting down...)' : ''}
</div>
</Popover>
{/if}