feat: kafka triggers (#4713)

* feat: kafka triggers

* sqlx

* fix build

* improve error messages on windows

* nit

* fix build

* missing action

* nit

* maybe fix ssl on windows

* nit

* update ee ref

---------

Co-authored-by: Ruben Fiszel <ruben@rubenfiszel.com>
Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
This commit is contained in:
HugoCasa
2024-11-18 17:03:50 +01:00
committed by GitHub
co-authored by Ruben Fiszel Ruben Fiszel
parent 6f4a6d8a72
commit 2a3c0d55fd
51 changed files with 1736 additions and 92 deletions
+16 -9
View File
@@ -8,11 +8,6 @@ use std::collections::HashMap;
* LICENSE-AGPL for a copy of the license.
*/
#[cfg(feature = "parquet")]
use crate::{job_helpers_ee::{
get_random_file_name, get_s3_resource, get_workspace_s3_resource, upload_file_internal,
UploadFileResponse,
}, users::fetch_api_authed_from_permissioned_as};
use crate::{
db::{ApiAuthed, DB},
resources::get_resource_value_interpolated_internal,
@@ -22,6 +17,14 @@ use crate::{
webhook_util::{WebhookMessage, WebhookShared},
HTTP_CLIENT,
};
#[cfg(feature = "parquet")]
use crate::{
job_helpers_ee::{
get_random_file_name, get_s3_resource, get_workspace_s3_resource, upload_file_internal,
UploadFileResponse,
},
users::fetch_api_authed_from_permissioned_as,
};
use axum::{
extract::{Extension, Json, Path, Query},
response::IntoResponse,
@@ -1288,7 +1291,6 @@ async fn upload_s3_file_from_app() -> Result<()> {
));
}
#[cfg(feature = "parquet")]
#[derive(Debug, Deserialize, Clone)]
struct UploadFileToS3Query {
@@ -1353,9 +1355,14 @@ async fn upload_s3_file_from_app(
let (username, permissioned_as, email) =
get_on_behalf_details_from_policy_and_authed(&policy, &opt_authed).await?;
let on_behalf_authed =
fetch_api_authed_from_permissioned_as(permissioned_as, email, &w_id, &db, username)
.await?;
let on_behalf_authed = fetch_api_authed_from_permissioned_as(
permissioned_as,
email,
&w_id,
&db,
Some(username),
)
.await?;
if let Some(file_key) = query.file_key {
// file key is provided => requires workspace, user or list policy and must match the regex
+3 -1
View File
@@ -23,7 +23,7 @@ use windmill_common::{
utils::{not_found_if_none, StripPath},
};
const KINDS: [&str; 10] = [
const KINDS: [&str; 12] = [
"script",
"group_",
"resource",
@@ -34,6 +34,8 @@ const KINDS: [&str; 10] = [
"app",
"raw_app",
"http_trigger",
"websocket_trigger",
"kafka_trigger",
];
pub fn workspaced_service() -> Router {
+1 -1
View File
@@ -507,7 +507,7 @@ async fn get_http_route_trigger(
trigger.email.clone(),
&trigger.workspace_id,
&db,
username_override.unwrap_or("anonymous".to_string()),
Some(username_override.unwrap_or("anonymous".to_string())),
)
.await?;
@@ -0,0 +1,14 @@
use crate::db::DB;
use axum::Router;
pub fn workspaced_service() -> Router {
Router::new()
}
pub async fn start_kafka_consumers(
_db: DB,
_rsmq: Option<rsmq_async::MultiplexedRsmq>,
mut _killpill_rx: tokio::sync::broadcast::Receiver<()>,
) -> () {
// implementation is not open source
}
+27 -6
View File
@@ -68,6 +68,8 @@ mod ai;
mod job_helpers_ee;
pub mod job_metrics;
pub mod jobs;
#[cfg(all(feature = "enterprise", feature = "kafka"))]
mod kafka_triggers_ee;
pub mod oauth2_ee;
mod oidc_ee;
mod raw_apps;
@@ -235,6 +237,9 @@ pub async fn run_server(
}
}
// #[cfg(feature = "kafka")]
// start_listening().await;
let job_helpers_service = {
#[cfg(feature = "parquet")]
{
@@ -247,9 +252,27 @@ pub async fn run_server(
}
};
let kafka_triggers_service = {
#[cfg(all(feature = "enterprise", feature = "kafka"))]
{
kafka_triggers_ee::workspaced_service()
}
#[cfg(not(all(feature = "enterprise", feature = "kafka")))]
{
Router::new()
}
};
if !*CLOUD_HOSTED {
let ws_killpill_rx = rx.resubscribe();
websocket_triggers::start_websockets(db.clone(), rsmq, ws_killpill_rx).await;
websocket_triggers::start_websockets(db.clone(), rsmq.clone(), ws_killpill_rx).await;
#[cfg(all(feature = "enterprise", feature = "kafka"))]
{
let kafka_killpill_rx = rx.resubscribe();
kafka_triggers_ee::start_kafka_consumers(db.clone(), rsmq, kafka_killpill_rx).await;
}
}
// build our application with a route
@@ -296,7 +319,8 @@ pub async fn run_server(
.nest(
"/websocket_triggers",
websocket_triggers::workspaced_service(),
),
)
.nest("/kafka_triggers", kafka_triggers_service),
)
.nest("/workspaces", workspaces::global_service())
.nest(
@@ -321,10 +345,7 @@ pub async fn run_server(
"/srch/w/:workspace_id/index",
indexer_ee::workspaced_service(),
)
.nest(
"/srch/index",
indexer_ee::global_service(),
)
.nest("/srch/index", indexer_ee::global_service())
.nest("/oidc", oidc_ee::global_service())
.nest(
"/saml",
+12
View File
@@ -18,6 +18,7 @@ pub struct TriggersCount {
webhook_count: i64,
email_count: i64,
websocket_count: i64,
kafka_count: i64,
}
pub(crate) async fn get_triggers_count_internal(
db: &DB,
@@ -64,6 +65,16 @@ pub(crate) async fn get_triggers_count_internal(
.await?
.unwrap_or(0);
let kafka_count = sqlx::query_scalar!(
"SELECT COUNT(*) FROM kafka_trigger WHERE script_path = $1 AND is_flow = $2 AND workspace_id = $3",
path,
is_flow,
w_id
)
.fetch_one(db)
.await?
.unwrap_or(0);
let webhook_count = (if is_flow {
sqlx::query_scalar!(
"SELECT COUNT(*) FROM token WHERE label LIKE 'webhook-%' AND workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
@@ -105,6 +116,7 @@ pub(crate) async fn get_triggers_count_internal(
webhook_count,
email_count,
websocket_count,
kafka_count,
}))
}
+3 -3
View File
@@ -746,7 +746,7 @@ pub async fn fetch_api_authed(
email: String,
w_id: &str,
db: &DB,
username_override: String,
username_override: Option<String>,
) -> error::Result<ApiAuthed> {
let permissioned_as = username_to_permissioned_as(username.as_str());
fetch_api_authed_from_permissioned_as(permissioned_as, email, w_id, db, username_override).await
@@ -757,7 +757,7 @@ pub async fn fetch_api_authed_from_permissioned_as(
email: String,
w_id: &str,
db: &DB,
username_override: String,
username_override: Option<String>,
) -> error::Result<ApiAuthed> {
let authed =
fetch_authed_from_permissioned_as(permissioned_as, email.clone(), w_id, db).await?;
@@ -769,7 +769,7 @@ pub async fn fetch_api_authed_from_permissioned_as(
groups: authed.groups,
folders: authed.folders,
scopes: authed.scopes,
username_override: Some(username_override),
username_override: username_override,
})
}
@@ -527,7 +527,7 @@ async fn wait_runnable_result(
ws_trigger.email.clone(),
&ws_trigger.workspace_id,
&db,
username_override,
Some(username_override),
)
.await?;
@@ -773,7 +773,9 @@ async fn listen_to_websocket(
rsmq: Option<rsmq_async::MultiplexedRsmq>,
mut killpill_rx: tokio::sync::broadcast::Receiver<()>,
) -> () {
update_ping(&db, &ws_trigger, Some("Connecting...")).await;
if let None = update_ping(&db, &ws_trigger, Some("Connecting...")).await {
return;
}
let url = ws_trigger.url.as_str();
@@ -970,7 +972,7 @@ async fn run_job(
trigger.email.clone(),
&trigger.workspace_id,
db,
"anonymous".to_string(),
Some("anonymous".to_string()),
)
.await?;
+5 -1
View File
@@ -1296,6 +1296,7 @@ async fn set_encryption_key(
struct UsedTriggers {
pub websocket_used: bool,
pub http_routes_used: bool,
pub kafka_used: bool,
}
async fn get_used_triggers(
@@ -1306,7 +1307,10 @@ async fn get_used_triggers(
let mut tx = user_db.begin(&authed).await?;
let websocket_used = sqlx::query_as!(
UsedTriggers,
r#"SELECT EXISTS(SELECT 1 FROM websocket_trigger WHERE workspace_id = $1) as "websocket_used!", EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1) as "http_routes_used!""#,
r#"SELECT
EXISTS(SELECT 1 FROM websocket_trigger WHERE workspace_id = $1) as "websocket_used!",
EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1) as "http_routes_used!",
EXISTS(SELECT 1 FROM kafka_trigger WHERE workspace_id = $1) as "kafka_used!""#,
w_id,
)
.fetch_one(&mut *tx)