diff --git a/backend/.sqlx/query-038d2fde90fa9e99e30d15161777fa3ab402e33edfca46daa95b52e525424586.json b/backend/.sqlx/query-038d2fde90fa9e99e30d15161777fa3ab402e33edfca46daa95b52e525424586.json new file mode 100644 index 0000000000..5dcf8e0792 --- /dev/null +++ b/backend/.sqlx/query-038d2fde90fa9e99e30d15161777fa3ab402e33edfca46daa95b52e525424586.json @@ -0,0 +1,41 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT id, topic, partition, \"offset\" FROM kafka_pending_commits\n WHERE workspace_id = $1 AND kafka_trigger_path = $2\n ORDER BY id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "topic", + "type_info": "Varchar" + }, + { + "ordinal": 2, + "name": "partition", + "type_info": "Int4" + }, + { + "ordinal": 3, + "name": "offset", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + false, + false, + false, + false + ] + }, + "hash": "038d2fde90fa9e99e30d15161777fa3ab402e33edfca46daa95b52e525424586" +} diff --git a/backend/.sqlx/query-12631fecee6aa11a45cf5c8d101c0dd8de50ac9b57e68198f637038d344fdd46.json b/backend/.sqlx/query-072e5ab78f929c6b7264f98c1588cb24cc635836276ee6faa2438f494bfbce04.json similarity index 54% rename from backend/.sqlx/query-12631fecee6aa11a45cf5c8d101c0dd8de50ac9b57e68198f637038d344fdd46.json rename to backend/.sqlx/query-072e5ab78f929c6b7264f98c1588cb24cc635836276ee6faa2438f494bfbce04.json index c452c33018..812c323e74 100644 --- a/backend/.sqlx/query-12631fecee6aa11a45cf5c8d101c0dd8de50ac9b57e68198f637038d344fdd46.json +++ b/backend/.sqlx/query-072e5ab78f929c6b7264f98c1588cb24cc635836276ee6faa2438f494bfbce04.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n UPDATE kafka_trigger\n SET\n kafka_resource_path = $1,\n group_id = $2,\n topics = $3,\n filters = $4,\n auto_offset_reset = $5,\n script_path = $6,\n path = $7,\n is_flow = $8,\n edited_by = $9,\n email = $10,\n edited_at = now(),\n server_id = NULL,\n error = NULL,\n error_handler_path = $13,\n error_handler_args = $14,\n retry = $15\n WHERE\n workspace_id = $11 AND path = $12\n ", + "query": "\n UPDATE kafka_trigger\n SET\n kafka_resource_path = $1,\n group_id = $2,\n topics = $3,\n filters = $4,\n auto_offset_reset = $5,\n auto_commit = $6,\n script_path = $7,\n path = $8,\n is_flow = $9,\n edited_by = $10,\n email = $11,\n edited_at = now(),\n server_id = NULL,\n error = NULL,\n error_handler_path = $14,\n error_handler_args = $15,\n retry = $16\n WHERE\n workspace_id = $12 AND path = $13\n ", "describe": { "columns": [], "parameters": { @@ -10,6 +10,7 @@ "VarcharArray", "JsonbArray", "Varchar", + "Bool", "Varchar", "Varchar", "Bool", @@ -24,5 +25,5 @@ }, "nullable": [] }, - "hash": "12631fecee6aa11a45cf5c8d101c0dd8de50ac9b57e68198f637038d344fdd46" + "hash": "072e5ab78f929c6b7264f98c1588cb24cc635836276ee6faa2438f494bfbce04" } diff --git a/backend/.sqlx/query-1df610a583e86edb70c374fd66c68554a6a4291426c09dd5b04fd832f9d31208.json b/backend/.sqlx/query-1df610a583e86edb70c374fd66c68554a6a4291426c09dd5b04fd832f9d31208.json new file mode 100644 index 0000000000..babf3ffbcc --- /dev/null +++ b/backend/.sqlx/query-1df610a583e86edb70c374fd66c68554a6a4291426c09dd5b04fd832f9d31208.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT reset_offset FROM kafka_trigger WHERE workspace_id = $1 AND path = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "reset_offset", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "1df610a583e86edb70c374fd66c68554a6a4291426c09dd5b04fd832f9d31208" +} diff --git a/backend/.sqlx/query-45fc21026fa76e5d69f00a68a7be81abb3ec627578f2d14f0ce33896dc6ab4cf.json b/backend/.sqlx/query-45fc21026fa76e5d69f00a68a7be81abb3ec627578f2d14f0ce33896dc6ab4cf.json new file mode 100644 index 0000000000..b5873760fc --- /dev/null +++ b/backend/.sqlx/query-45fc21026fa76e5d69f00a68a7be81abb3ec627578f2d14f0ce33896dc6ab4cf.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO kafka_trigger (\n path, kafka_resource_path, topics, group_id, script_path,\n is_flow, workspace_id, edited_by, email, auto_commit\n )\n VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "VarcharArray", + "Varchar", + "Varchar", + "Bool", + "Varchar", + "Varchar", + "Varchar", + "Bool" + ] + }, + "nullable": [] + }, + "hash": "45fc21026fa76e5d69f00a68a7be81abb3ec627578f2d14f0ce33896dc6ab4cf" +} diff --git a/backend/.sqlx/query-4b2a29b3ef7ec4802d81ec4b706623b991c938e40d0db25290b03dc0577c2740.json b/backend/.sqlx/query-4b2a29b3ef7ec4802d81ec4b706623b991c938e40d0db25290b03dc0577c2740.json new file mode 100644 index 0000000000..6dd5c799e7 --- /dev/null +++ b/backend/.sqlx/query-4b2a29b3ef7ec4802d81ec4b706623b991c938e40d0db25290b03dc0577c2740.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT auto_commit FROM kafka_trigger WHERE workspace_id = $1 AND path = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "auto_commit", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "4b2a29b3ef7ec4802d81ec4b706623b991c938e40d0db25290b03dc0577c2740" +} diff --git a/backend/.sqlx/query-7e3bfb33fb771aec39b43a7550091ce7c9b1261b52d10f4a7f3273fed3c916df.json b/backend/.sqlx/query-4cf4be7a981173d3f242887d9313c7e60d23e9827f23c0de5b546ed56697d54a.json similarity index 61% rename from backend/.sqlx/query-7e3bfb33fb771aec39b43a7550091ce7c9b1261b52d10f4a7f3273fed3c916df.json rename to backend/.sqlx/query-4cf4be7a981173d3f242887d9313c7e60d23e9827f23c0de5b546ed56697d54a.json index 23652a2571..9b76b7f048 100644 --- a/backend/.sqlx/query-7e3bfb33fb771aec39b43a7550091ce7c9b1261b52d10f4a7f3273fed3c916df.json +++ b/backend/.sqlx/query-4cf4be7a981173d3f242887d9313c7e60d23e9827f23c0de5b546ed56697d54a.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT kafka_resource_path, topics, group_id, mode AS \"mode: String\"\n FROM kafka_trigger\n WHERE workspace_id = $1 AND path = $2\n ", + "query": "\n SELECT kafka_resource_path, topics, group_id, mode AS \"mode: String\",\n auto_offset_reset, auto_commit, reset_offset\n FROM kafka_trigger\n WHERE workspace_id = $1 AND path = $2\n ", "describe": { "columns": [ { @@ -33,6 +33,21 @@ } } } + }, + { + "ordinal": 4, + "name": "auto_offset_reset", + "type_info": "Varchar" + }, + { + "ordinal": 5, + "name": "auto_commit", + "type_info": "Bool" + }, + { + "ordinal": 6, + "name": "reset_offset", + "type_info": "Bool" } ], "parameters": { @@ -42,11 +57,14 @@ ] }, "nullable": [ + false, + false, + false, false, false, false, false ] }, - "hash": "7e3bfb33fb771aec39b43a7550091ce7c9b1261b52d10f4a7f3273fed3c916df" + "hash": "4cf4be7a981173d3f242887d9313c7e60d23e9827f23c0de5b546ed56697d54a" } diff --git a/backend/.sqlx/query-50807b807bb901a380926798be655c13a18dfd26e237a8218d3006e2898b5aa3.json b/backend/.sqlx/query-50807b807bb901a380926798be655c13a18dfd26e237a8218d3006e2898b5aa3.json new file mode 100644 index 0000000000..ec667519e4 --- /dev/null +++ b/backend/.sqlx/query-50807b807bb901a380926798be655c13a18dfd26e237a8218d3006e2898b5aa3.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT auto_commit\n FROM kafka_trigger\n WHERE workspace_id = $1 AND path = $2\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "auto_commit", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "50807b807bb901a380926798be655c13a18dfd26e237a8218d3006e2898b5aa3" +} diff --git a/backend/.sqlx/query-4b5a711986017654bdd495893a16ddc6ab09c98cd8723865cbd341404bc6a02f.json b/backend/.sqlx/query-5dd6315ec270c268e905262e4b0a920837354d91a0ae16b1236c1267da71765f.json similarity index 63% rename from backend/.sqlx/query-4b5a711986017654bdd495893a16ddc6ab09c98cd8723865cbd341404bc6a02f.json rename to backend/.sqlx/query-5dd6315ec270c268e905262e4b0a920837354d91a0ae16b1236c1267da71765f.json index d90a467380..936a4650f6 100644 --- a/backend/.sqlx/query-4b5a711986017654bdd495893a16ddc6ab09c98cd8723865cbd341404bc6a02f.json +++ b/backend/.sqlx/query-5dd6315ec270c268e905262e4b0a920837354d91a0ae16b1236c1267da71765f.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n INSERT INTO kafka_trigger (\n workspace_id,\n path,\n kafka_resource_path,\n group_id,\n topics,\n filters,\n auto_offset_reset,\n script_path,\n is_flow,\n mode,\n edited_by,\n email,\n edited_at,\n error_handler_path,\n error_handler_args,\n retry\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, now(), $13, $14, $15\n )\n ", + "query": "\n INSERT INTO kafka_trigger (\n workspace_id,\n path,\n kafka_resource_path,\n group_id,\n topics,\n filters,\n auto_offset_reset,\n auto_commit,\n script_path,\n is_flow,\n mode,\n edited_by,\n email,\n edited_at,\n error_handler_path,\n error_handler_args,\n retry\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, now(), $14, $15, $16\n )\n ", "describe": { "columns": [], "parameters": { @@ -12,6 +12,7 @@ "VarcharArray", "JsonbArray", "Varchar", + "Bool", "Varchar", "Bool", { @@ -35,5 +36,5 @@ }, "nullable": [] }, - "hash": "4b5a711986017654bdd495893a16ddc6ab09c98cd8723865cbd341404bc6a02f" + "hash": "5dd6315ec270c268e905262e4b0a920837354d91a0ae16b1236c1267da71765f" } diff --git a/backend/.sqlx/query-80bad96cbec6b5eca57a6380e7515565490a271050dcc4b5aac2b730ae3a55b9.json b/backend/.sqlx/query-80bad96cbec6b5eca57a6380e7515565490a271050dcc4b5aac2b730ae3a55b9.json new file mode 100644 index 0000000000..31f2767fe8 --- /dev/null +++ b/backend/.sqlx/query-80bad96cbec6b5eca57a6380e7515565490a271050dcc4b5aac2b730ae3a55b9.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM kafka_pending_commits WHERE id = ANY($1)", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Int8Array" + ] + }, + "nullable": [] + }, + "hash": "80bad96cbec6b5eca57a6380e7515565490a271050dcc4b5aac2b730ae3a55b9" +} diff --git a/backend/.sqlx/query-c2f38c9e09aac73d10e8f327715927c07832badb2c9145d5996b829163bdf7d9.json b/backend/.sqlx/query-c2f38c9e09aac73d10e8f327715927c07832badb2c9145d5996b829163bdf7d9.json new file mode 100644 index 0000000000..7415b311f1 --- /dev/null +++ b/backend/.sqlx/query-c2f38c9e09aac73d10e8f327715927c07832badb2c9145d5996b829163bdf7d9.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE kafka_trigger SET reset_offset = true, server_id = NULL WHERE workspace_id = $1 AND path = $2 RETURNING true", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "?column?", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "c2f38c9e09aac73d10e8f327715927c07832badb2c9145d5996b829163bdf7d9" +} diff --git a/backend/.sqlx/query-ef15599f532fab2cbb487542ffec047cf3b7ce22ce868db1b1a63e6c10d0d12b.json b/backend/.sqlx/query-ef15599f532fab2cbb487542ffec047cf3b7ce22ce868db1b1a63e6c10d0d12b.json new file mode 100644 index 0000000000..238977d563 --- /dev/null +++ b/backend/.sqlx/query-ef15599f532fab2cbb487542ffec047cf3b7ce22ce868db1b1a63e6c10d0d12b.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE kafka_trigger SET reset_offset = false WHERE workspace_id = $1 AND path = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "ef15599f532fab2cbb487542ffec047cf3b7ce22ce868db1b1a63e6c10d0d12b" +} diff --git a/backend/.sqlx/query-f67e5c96eb9cb35953d4c3e83e0fcbb5b647737e0366529a2f418218b1a74679.json b/backend/.sqlx/query-f67e5c96eb9cb35953d4c3e83e0fcbb5b647737e0366529a2f418218b1a74679.json new file mode 100644 index 0000000000..5797d761e9 --- /dev/null +++ b/backend/.sqlx/query-f67e5c96eb9cb35953d4c3e83e0fcbb5b647737e0366529a2f418218b1a74679.json @@ -0,0 +1,18 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO kafka_pending_commits (workspace_id, kafka_trigger_path, topic, partition, \"offset\")\n VALUES ($1, $2, $3, $4, $5)", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Varchar", + "Int4", + "Int8" + ] + }, + "nullable": [] + }, + "hash": "f67e5c96eb9cb35953d4c3e83e0fcbb5b647737e0366529a2f418218b1a74679" +} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 23e6f2a4b1..b10179c658 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -16870,7 +16870,7 @@ dependencies = [ [[package]] name = "windmill-parser-py-asset" -version = "1.653.0" +version = "1.654.0" dependencies = [ "anyhow", "rustpython-ast", @@ -16950,7 +16950,7 @@ dependencies = [ [[package]] name = "windmill-parser-sql-asset" -version = "1.653.0" +version = "1.654.0" dependencies = [ "anyhow", "serde", @@ -16980,7 +16980,7 @@ dependencies = [ [[package]] name = "windmill-parser-ts-asset" -version = "1.653.0" +version = "1.654.0" dependencies = [ "anyhow", "serde-wasm-bindgen", diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 18e2279e37..a048a9b463 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -85bcbfc8e952842c555156d1050b2a024e7e37a3 \ No newline at end of file +fcd3ea52b0cc94fbe1159baf662a38da947456de \ No newline at end of file diff --git a/backend/migrations/20260312000000_kafka_auto_commit.down.sql b/backend/migrations/20260312000000_kafka_auto_commit.down.sql new file mode 100644 index 0000000000..58d1ad2659 --- /dev/null +++ b/backend/migrations/20260312000000_kafka_auto_commit.down.sql @@ -0,0 +1,2 @@ +DROP TABLE IF EXISTS kafka_pending_commits; +ALTER TABLE kafka_trigger DROP COLUMN auto_commit; diff --git a/backend/migrations/20260312000000_kafka_auto_commit.up.sql b/backend/migrations/20260312000000_kafka_auto_commit.up.sql new file mode 100644 index 0000000000..643e2ce1f3 --- /dev/null +++ b/backend/migrations/20260312000000_kafka_auto_commit.up.sql @@ -0,0 +1,14 @@ +ALTER TABLE kafka_trigger ADD COLUMN auto_commit BOOLEAN NOT NULL DEFAULT TRUE; + +CREATE TABLE kafka_pending_commits ( + id BIGSERIAL PRIMARY KEY, + workspace_id VARCHAR(50) NOT NULL, + kafka_trigger_path VARCHAR(255) NOT NULL, + topic VARCHAR(255) NOT NULL, + partition INTEGER NOT NULL, + "offset" BIGINT NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + FOREIGN KEY (workspace_id, kafka_trigger_path) REFERENCES kafka_trigger(workspace_id, path) ON DELETE CASCADE +); + +CREATE INDEX idx_kafka_pending_commits_trigger ON kafka_pending_commits (workspace_id, kafka_trigger_path); diff --git a/backend/summarized_schema.txt b/backend/summarized_schema.txt index 3ece1ce7ef..2f86dcb322 100644 --- a/backend/summarized_schema.txt +++ b/backend/summarized_schema.txt @@ -109,7 +109,9 @@ job_result_stream_v2: job_id(uuid), workspace_id(text), stream(text), idx(int) job_settings: job_id(uuid), runnable_settings(bigint) job_stats: workspace_id(char), job_id(uuid), metric_id(char), metric_name(char), metric_kind(metric_kind), scalar_int(int), scalar_float(float), timestamps(ts), timeseries_int(int[]), timeseries_float(float[]) FK: (workspace_id) -> workspace(id) -kafka_trigger: path(char), kafka_resource_path(char), topics(char), group_id(char), script_path(char), is_flow(bool), workspace_id(char), edited_by(char), email(char), edited_at(ts), extra_perms(jsonb), server_id(char), last_server_ping(ts), error(text), error_handler_path(char), error_handler_args(jsonb), retry(jsonb), mode(trigger_mode), filters(jsonb[]) +kafka_pending_commits: id(bigint), workspace_id(char), kafka_trigger_path(char), topic(char), partition(int), offset(bigint), created_at(ts) + FK: (workspace_id, kafka_trigger_path) -> kafka_trigger(workspace_id, path) +kafka_trigger: path(char), kafka_resource_path(char), topics(char), group_id(char), script_path(char), is_flow(bool), workspace_id(char), edited_by(char), email(char), edited_at(ts), extra_perms(jsonb), server_id(char), last_server_ping(ts), error(text), error_handler_path(char), error_handler_args(jsonb), retry(jsonb), mode(trigger_mode), filters(jsonb[]), auto_commit(bool) log_file: hostname(char), log_ts(ts), ok_lines(bigint), err_lines(bigint), mode(log_mode), worker_group(char), file_path(char), json_fmt(bool) magic_link: email(char), token(char), expiration(ts) mcp_oauth_client: mcp_server_url(text), client_id(text), client_secret(text), client_secret_expires_at(ts), token_endpoint(text), created_at(ts) diff --git a/backend/windmill-api-integration-tests/tests/trigger_e2e.rs b/backend/windmill-api-integration-tests/tests/trigger_e2e.rs index 3b0e32749a..1d4f122bcc 100644 --- a/backend/windmill-api-integration-tests/tests/trigger_e2e.rs +++ b/backend/windmill-api-integration-tests/tests/trigger_e2e.rs @@ -251,7 +251,8 @@ async fn test_websocket_e2e(db: Pool) -> anyhow::Result<()> { "test-workspace", "test-user", "test@windmill.dev", - &[json!({"type": "RawMessage", "content": "hello from e2e test"})] as &[serde_json::Value], + &[json!({"type": "RawMessage", "content": "hello from e2e test"})] + as &[serde_json::Value], ) .execute(&db) .await?; @@ -302,9 +303,11 @@ async fn test_postgres_e2e(db: Pool) -> anyhow::Result<()> { sqlx::query("CREATE TABLE test_trigger_table (id serial PRIMARY KEY, data text)") .execute(&db) .await?; - sqlx::query(&format!("CREATE PUBLICATION {pub_name} FOR TABLE test_trigger_table")) - .execute(&db) - .await?; + sqlx::query(&format!( + "CREATE PUBLICATION {pub_name} FOR TABLE test_trigger_table" + )) + .execute(&db) + .await?; sqlx::query(&format!( "SELECT pg_create_logical_replication_slot('{slot_name}', 'pgoutput')" )) @@ -313,10 +316,9 @@ async fn test_postgres_e2e(db: Pool) -> anyhow::Result<()> { // Extract the test DB name from the pool so the resource points here, // not at the main windmill database. - let test_db_name: String = - sqlx::query_scalar("SELECT current_database()") - .fetch_one(&db) - .await?; + let test_db_name: String = sqlx::query_scalar("SELECT current_database()") + .fetch_one(&db) + .await?; insert_resource( &db, diff --git a/backend/windmill-api-integration-tests/tests/triggers.rs b/backend/windmill-api-integration-tests/tests/triggers.rs index 46854d566e..b55f24397c 100644 --- a/backend/windmill-api-integration-tests/tests/triggers.rs +++ b/backend/windmill-api-integration-tests/tests/triggers.rs @@ -214,12 +214,9 @@ async fn test_capture_delete(db: Pool) -> anyhow::Result<()> { .execute(&db) .await?; - let count = sqlx::query_scalar!( - "SELECT COUNT(*) FROM capture WHERE id = $1", - id, - ) - .fetch_one(&db) - .await?; + let count = sqlx::query_scalar!("SELECT COUNT(*) FROM capture WHERE id = $1", id,) + .fetch_one(&db) + .await?; assert_eq!(count, Some(0)); @@ -385,7 +382,10 @@ async fn test_capture_api_list_captures(db: Pool) -> anyhow::Result<() .send() .await?; - assert!(response.status().is_success(), "list captures should succeed"); + assert!( + response.status().is_success(), + "list captures should succeed" + ); let captures: Vec = response.json().await?; assert_eq!(captures.len(), 3); @@ -480,12 +480,9 @@ async fn test_capture_api_delete(db: Pool) -> anyhow::Result<()> { assert!(response.status().is_success(), "delete should succeed"); - let count = sqlx::query_scalar!( - "SELECT COUNT(*) FROM capture WHERE id = $1", - id, - ) - .fetch_one(&db) - .await?; + let count = sqlx::query_scalar!("SELECT COUNT(*) FROM capture WHERE id = $1", id,) + .fetch_one(&db) + .await?; assert_eq!(count, Some(0)); @@ -933,7 +930,8 @@ async fn test_kafka_trigger_insert(db: Pool) -> anyhow::Result<()> { let trigger = sqlx::query!( r#" - SELECT kafka_resource_path, topics, group_id, mode AS "mode: String" + SELECT kafka_resource_path, topics, group_id, mode AS "mode: String", + auto_offset_reset, auto_commit, reset_offset FROM kafka_trigger WHERE workspace_id = $1 AND path = $2 "#, @@ -947,6 +945,50 @@ async fn test_kafka_trigger_insert(db: Pool) -> anyhow::Result<()> { assert_eq!(trigger.topics, vec!["topic-a", "topic-b"]); assert_eq!(trigger.group_id, "my-consumer-group"); assert_eq!(trigger.mode, "enabled"); + assert_eq!(trigger.auto_offset_reset, "latest"); + assert_eq!(trigger.auto_commit, true); + assert_eq!(trigger.reset_offset, false); + + Ok(()) +} + +#[sqlx::test(migrations = "../migrations", fixtures("base"))] +async fn test_kafka_trigger_insert_auto_commit_disabled(db: Pool) -> anyhow::Result<()> { + sqlx::query!( + r#" + INSERT INTO kafka_trigger ( + path, kafka_resource_path, topics, group_id, script_path, + is_flow, workspace_id, edited_by, email, auto_commit + ) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) + "#, + "f/test/kafka_trigger_no_commit", + "u/admin/kafka_resource", + &["topic-c"] as &[&str], + "my-consumer-group-2", + "f/test/kafka_handler", + false, + "test-workspace", + "test-user", + "test@windmill.dev", + false, + ) + .execute(&db) + .await?; + + let trigger = sqlx::query!( + r#" + SELECT auto_commit + FROM kafka_trigger + WHERE workspace_id = $1 AND path = $2 + "#, + "test-workspace", + "f/test/kafka_trigger_no_commit", + ) + .fetch_one(&db) + .await?; + + assert_eq!(trigger.auto_commit, false); Ok(()) } diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 9bb25f5589..9ec79f84d4 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -12017,6 +12017,39 @@ paths: "200": description: kafka trigger offsets reset successfully + /w/{workspace}/kafka_triggers/commit_offsets/{path}: + post: + summary: commit kafka offsets for a specific trigger + operationId: commitKafkaOffsets + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Path" + requestBody: + description: offsets to commit + required: true + content: + application/json: + schema: + type: object + properties: + topic: + type: string + partition: + type: integer + format: int32 + offset: + type: integer + format: int64 + required: + - topic + - partition + - offset + responses: + "200": + description: kafka offsets committed successfully + /w/{workspace}/nats_triggers/create: post: summary: create nats trigger @@ -22098,6 +22131,10 @@ components: - earliest default: latest description: "Initial offset behavior when consumer group has no committed offset. 'latest' starts from new messages only, 'earliest' starts from the beginning." + auto_commit: + type: boolean + default: true + description: "When true (default), offsets are committed automatically after receiving each message. When false, you must manually commit offsets using the commit_offsets endpoint." server_id: type: string description: ID of the server currently handling this trigger (internal) @@ -22165,6 +22202,10 @@ components: - earliest default: latest description: "Initial offset behavior when consumer group has no committed offset." + auto_commit: + type: boolean + default: true + description: "When true (default), offsets are committed automatically after receiving each message. When false, you must manually commit offsets using the commit_offsets endpoint." mode: $ref: "#/components/schemas/TriggerMode" error_handler_path: @@ -22224,6 +22265,10 @@ components: - earliest default: latest description: "Initial offset behavior when consumer group has no committed offset." + auto_commit: + type: boolean + default: true + description: "When true (default), offsets are committed automatically after receiving each message. When false, you must manually commit offsets using the commit_offsets endpoint." path: type: string description: The unique path identifier for this trigger diff --git a/backend/windmill-api/src/capture.rs b/backend/windmill-api/src/capture.rs index bcc287ec04..277504b481 100644 --- a/backend/windmill-api/src/capture.rs +++ b/backend/windmill-api/src/capture.rs @@ -1012,6 +1012,7 @@ async fn http_payload( .to_v2_preprocessor_args( &http_trigger_config.route_path, &route_path, + "", ¶ms, headers, query, diff --git a/backend/windmill-api/src/triggers/http/handler.rs b/backend/windmill-api/src/triggers/http/handler.rs index 9d95ead840..fcb883f089 100644 --- a/backend/windmill-api/src/triggers/http/handler.rs +++ b/backend/windmill-api/src/triggers/http/handler.rs @@ -468,6 +468,7 @@ async fn route_job( .to_args_from_format( &trigger.route_path, &called_path, + &trigger.path, ¶ms, runnable_format, trigger.wrap_body, diff --git a/backend/windmill-api/src/triggers/http/http_trigger_args.rs b/backend/windmill-api/src/triggers/http/http_trigger_args.rs index 4866f9be34..7df1e9f200 100644 --- a/backend/windmill-api/src/triggers/http/http_trigger_args.rs +++ b/backend/windmill-api/src/triggers/http/http_trigger_args.rs @@ -68,6 +68,7 @@ struct HttpTriggerPreprocessorEvent<'a> { kind: String, route: &'a str, path: &'a str, + trigger_path: &'a str, body: Box, raw_string: Option, params: &'a HashMap, @@ -117,6 +118,7 @@ impl HttpTriggerArgs { self, route_path: &str, called_path: &str, + trigger_path: &str, params: &HashMap, format: RunnableFormat, wrap_body: bool, @@ -126,7 +128,14 @@ impl HttpTriggerArgs { match format { RunnableFormat { has_preprocessor: true, version: RunnableFormatVersion::V2 } => { // we don't care about wrap_body in v2 - self.to_v2_preprocessor_args(route_path, called_path, params, headers, query) + self.to_v2_preprocessor_args( + route_path, + called_path, + trigger_path, + params, + headers, + query, + ) } RunnableFormat { has_preprocessor: true, version: RunnableFormatVersion::V1 } => self .to_v1_preprocessor_args( @@ -177,6 +186,7 @@ impl HttpTriggerArgs { self, route_path: &str, called_path: &str, + trigger_path: &str, params: &HashMap, headers: HashMap>, query: HashMap>, @@ -193,6 +203,7 @@ impl HttpTriggerArgs { method: (&self.0.metadata.method).try_into()?, route: route_path, path: called_path, + trigger_path, params, }), ); diff --git a/backend/windmill-trigger-websocket/src/listener.rs b/backend/windmill-trigger-websocket/src/listener.rs index f6fba882a8..8aab61d9c6 100644 --- a/backend/windmill-trigger-websocket/src/listener.rs +++ b/backend/windmill-trigger-websocket/src/listener.rs @@ -322,7 +322,7 @@ impl Listener for WebsocketTrigger { db: &DB, listening_trigger: &ListeningTrigger, payload: Self::Payload, - trigger_info: HashMap>, + mut trigger_info: HashMap>, extra: Option, ) -> Result<()> { let ListeningTrigger { @@ -338,6 +338,7 @@ impl Listener for WebsocketTrigger { let WebsocketConfig { url, .. } = trigger_config; + trigger_info.insert("trigger_path".to_string(), to_raw_value(path)); let args = WebsocketTrigger::build_job_args( &script_path, *is_flow, diff --git a/backend/windmill-trigger/src/listener.rs b/backend/windmill-trigger/src/listener.rs index d09bb5e55a..04e1df102d 100644 --- a/backend/windmill-trigger/src/listener.rs +++ b/backend/windmill-trigger/src/listener.rs @@ -21,6 +21,7 @@ use windmill_common::{ jobs::JobTriggerKind, triggers::{TriggerKind, TriggerMetadata}, utils::report_critical_error, + worker::to_raw_value, DB, INSTANCE_NAME, }; @@ -467,9 +468,13 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { db: &DB, listening_trigger: &ListeningTrigger, payload: Self::Payload, - trigger_info: HashMap>, + mut trigger_info: HashMap>, _extra: Option, ) -> Result<()> { + trigger_info.insert( + "trigger_path".to_string(), + to_raw_value(&listening_trigger.path), + ); let args = Self::build_job_args( &listening_trigger.script_path, listening_trigger.is_flow, @@ -552,6 +557,11 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { return Ok(()); } + let mut trigger_info = trigger_info; + trigger_info.insert( + "trigger_path".to_string(), + to_raw_value(&listening_trigger.path), + ); let (main_args, preprocessor_args) = Self::build_capture_payloads(&payload, trigger_info); if let Err(err) = insert_capture_payload( db, diff --git a/frontend/src/lib/components/triggers/TriggerAdvancedBadges.svelte b/frontend/src/lib/components/triggers/TriggerAdvancedBadges.svelte new file mode 100644 index 0000000000..b261aa8ed7 --- /dev/null +++ b/frontend/src/lib/components/triggers/TriggerAdvancedBadges.svelte @@ -0,0 +1,28 @@ + + +{#if allBadges.length > 0} +
+ {#each allBadges as badge} + {badge.name} + {/each} +
+{/if} diff --git a/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte index b40ef8e505..a33f14243b 100644 --- a/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte @@ -25,6 +25,7 @@ import Tabs from '$lib/components/common/tabs/Tabs.svelte' import Tab from '$lib/components/common/tabs/Tab.svelte' import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte' + import TriggerAdvancedBadges from '../TriggerAdvancedBadges.svelte' import { saveEmailTriggerFromCfg } from './utils' import { deepEqual } from 'fast-equals' import TriggerSuspendedJobsAlert from '../TriggerSuspendedJobsAlert.svelte' @@ -370,7 +371,10 @@ />
-
+ {#snippet header()} + + {/snippet} +
@@ -390,6 +394,7 @@
+
{/if} {/snippet} diff --git a/frontend/src/lib/components/triggers/gcp/GcpTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/gcp/GcpTriggerEditorInner.svelte index 5e1b08d2c3..6a11d6180b 100644 --- a/frontend/src/lib/components/triggers/gcp/GcpTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/gcp/GcpTriggerEditorInner.svelte @@ -31,6 +31,7 @@ import Tabs from '$lib/components/common/tabs/Tabs.svelte' import Tab from '$lib/components/common/tabs/Tab.svelte' import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte' + import TriggerAdvancedBadges from '../TriggerAdvancedBadges.svelte' import Subsection from '$lib/components/Subsection.svelte' import Toggle from '$lib/components/Toggle.svelte' @@ -452,7 +453,12 @@ />
-
+ {#snippet header()} + + {/snippet} +
@@ -520,6 +526,7 @@
+
{/if} {/snippet} diff --git a/frontend/src/lib/components/triggers/http/RouteEditorInner.svelte b/frontend/src/lib/components/triggers/http/RouteEditorInner.svelte index a35227e081..272a0b707c 100644 --- a/frontend/src/lib/components/triggers/http/RouteEditorInner.svelte +++ b/frontend/src/lib/components/triggers/http/RouteEditorInner.svelte @@ -52,6 +52,7 @@ import Tabs from '$lib/components/common/tabs/Tabs.svelte' import Tab from '$lib/components/common/tabs/Tab.svelte' import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte' + import TriggerAdvancedBadges from '../TriggerAdvancedBadges.svelte' import { deepEqual } from 'fast-equals' import TriggerSuspendedJobsAlert from '../TriggerSuspendedJobsAlert.svelte' import TriggerSuspendedJobsModal from '../TriggerSuspendedJobsModal.svelte' @@ -695,6 +696,13 @@ {#if !is_static_website}
+ {#snippet header()} + + {/snippet}
@@ -909,6 +917,7 @@
+
{/if}
{/if} diff --git a/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte index ad5cce6f92..5f68f03705 100644 --- a/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/kafka/KafkaTriggerEditorInner.svelte @@ -6,12 +6,7 @@ import Path from '$lib/components/Path.svelte' import Required from '$lib/components/Required.svelte' import ScriptPicker from '$lib/components/ScriptPicker.svelte' - import { - KafkaTriggerService, - type ErrorHandler, - type Retry, - type TriggerMode - } from '$lib/gen' + import { KafkaTriggerService, type ErrorHandler, type Retry, type TriggerMode } from '$lib/gen' import { usedTriggerKinds, userStore, workspaceStore } from '$lib/stores' import { canWrite, capitalize, emptyString, sendUserToast } from '$lib/utils' import Section from '$lib/components/Section.svelte' @@ -25,10 +20,13 @@ import Tabs from '$lib/components/common/tabs/Tabs.svelte' import Tab from '$lib/components/common/tabs/Tab.svelte' import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte' + import TriggerAdvancedBadges from '../TriggerAdvancedBadges.svelte' import { deepEqual } from 'fast-equals' import TriggerSuspendedJobsAlert from '../TriggerSuspendedJobsAlert.svelte' import TriggerSuspendedJobsModal from '../TriggerSuspendedJobsModal.svelte' import TriggerFilters from '../TriggerFilters.svelte' + import Select from '$lib/components/select/Select.svelte' + import Toggle from '$lib/components/Toggle.svelte' interface Props { useDrawer?: boolean @@ -87,6 +85,7 @@ let kafkaResourcePath = $state('') let kafkaCfg: Record = $state({}) let autoOffsetReset = $state('latest') + let autoCommit = $state(true) let deploymentLoading = $state(false) let resetLoading = $state(false) let optionTabSelected: 'error_handler' | 'retries' = $state('error_handler') @@ -176,6 +175,7 @@ topics: nDefaultValues?.topics ?? [''] } autoOffsetReset = nDefaultValues?.auto_offset_reset ?? 'latest' + autoCommit = nDefaultValues?.auto_commit ?? true initialScriptPath = '' fixedScriptPath = fixedScriptPath_ ?? '' script_path = fixedScriptPath @@ -207,6 +207,7 @@ topics: cfg?.topics } autoOffsetReset = cfg?.auto_offset_reset ?? 'latest' + autoCommit = cfg?.auto_commit ?? true mode = cfg?.mode ?? 'enabled' extra_perms = cfg?.extra_perms can_write = canWrite(path, cfg?.extra_perms, $userStore) @@ -240,6 +241,7 @@ topics: kafkaCfg.topics, filters, auto_offset_reset: autoOffsetReset, + auto_commit: autoCommit, mode, extra_perms: extra_perms, error_handler_path, @@ -481,36 +483,85 @@ bind:kafkaCfgValid bind:kafkaResourcePath bind:kafkaCfg - bind:autoOffsetReset {path} {can_write} showTestingBadge={isEditor} /> - {#if edit && can_write} - - {/if} - - -
-
+ {#snippet header()} + 0 } + ]} + /> + {/snippet} +
+ -
diff --git a/frontend/src/lib/components/triggers/kafka/utils.ts b/frontend/src/lib/components/triggers/kafka/utils.ts index 16e065093d..12ee49ac24 100644 --- a/frontend/src/lib/components/triggers/kafka/utils.ts +++ b/frontend/src/lib/components/triggers/kafka/utils.ts @@ -25,6 +25,7 @@ export async function saveKafkaTriggerFromCfg( topics: cfg.topics, filters: cfg.filters ?? [], auto_offset_reset: cfg.auto_offset_reset ?? 'latest', + auto_commit: cfg.auto_commit ?? true, ...errorHandlerAndRetries } try { diff --git a/frontend/src/lib/components/triggers/mqtt/MqttTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/mqtt/MqttTriggerEditorInner.svelte index 18148fdbab..df5c7c21e1 100644 --- a/frontend/src/lib/components/triggers/mqtt/MqttTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/mqtt/MqttTriggerEditorInner.svelte @@ -29,6 +29,7 @@ import Tabs from '$lib/components/common/tabs/Tabs.svelte' import Tab from '$lib/components/common/tabs/Tab.svelte' import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte' + import TriggerAdvancedBadges from '../TriggerAdvancedBadges.svelte' import Toggle from '$lib/components/Toggle.svelte' import ToggleButtonGroup from '$lib/components/common/toggleButton-v2/ToggleButtonGroup.svelte' import ToggleButton from '$lib/components/common/toggleButton-v2/ToggleButton.svelte' @@ -469,7 +470,13 @@ />
-
+ {#snippet header()} + + {/snippet} +
@@ -606,6 +613,7 @@
+
{/if} {/snippet} diff --git a/frontend/src/lib/components/triggers/nats/NatsTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/nats/NatsTriggerEditorInner.svelte index 8c7fbb846c..a3c02c57dc 100644 --- a/frontend/src/lib/components/triggers/nats/NatsTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/nats/NatsTriggerEditorInner.svelte @@ -19,6 +19,7 @@ import Tabs from '$lib/components/common/tabs/Tabs.svelte' import Tab from '$lib/components/common/tabs/Tab.svelte' import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte' + import TriggerAdvancedBadges from '../TriggerAdvancedBadges.svelte' import { deepEqual } from 'fast-equals' import TriggerSuspendedJobsAlert from '../TriggerSuspendedJobsAlert.svelte' import TriggerSuspendedJobsModal from '../TriggerSuspendedJobsModal.svelte' @@ -452,7 +453,10 @@ />
-
+ {#snippet header()} + + {/snippet} +
@@ -472,6 +476,7 @@
+
{/if} {/snippet} diff --git a/frontend/src/lib/components/triggers/postgres/PostgresTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/postgres/PostgresTriggerEditorInner.svelte index cd176ac51c..5ffe489a22 100644 --- a/frontend/src/lib/components/triggers/postgres/PostgresTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/postgres/PostgresTriggerEditorInner.svelte @@ -36,6 +36,7 @@ import TestingBadge from '../testingBadge.svelte' import { getHandlerType, handleConfigChange, type Trigger } from '../utils' import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte' + import TriggerAdvancedBadges from '../TriggerAdvancedBadges.svelte' import { fade } from 'svelte/transition' import MultiSelect from '$lib/components/select/MultiSelect.svelte' import { safeSelectItems } from '$lib/components/select/utils.svelte' @@ -854,7 +855,10 @@
-
+ {#snippet header()} + + {/snippet} +
@@ -874,6 +878,7 @@
+
{/if} {/snippet} diff --git a/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte b/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte index 4aca94c4b5..55f825b40c 100644 --- a/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte +++ b/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte @@ -1,5 +1,6 @@