diff --git a/backend/.sqlx/query-023cdbc77ea9e2c17a1aa92a5b9001f29e58e81b3f782887db6e0a627dd8ad75.json b/backend/.sqlx/query-023cdbc77ea9e2c17a1aa92a5b9001f29e58e81b3f782887db6e0a627dd8ad75.json index d3f1c39c7a..5b58cedc72 100644 --- a/backend/.sqlx/query-023cdbc77ea9e2c17a1aa92a5b9001f29e58e81b3f782887db6e0a627dd8ad75.json +++ b/backend/.sqlx/query-023cdbc77ea9e2c17a1aa92a5b9001f29e58e81b3f782887db6e0a627dd8ad75.json @@ -12,8 +12,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-06072cfe26abe58629623a8b38382b33947c5a5c702ce586e6e6ea51430380bf.json b/backend/.sqlx/query-06072cfe26abe58629623a8b38382b33947c5a5c702ce586e6e6ea51430380bf.json index 97e39aa0e5..8ecc409fc8 100644 --- a/backend/.sqlx/query-06072cfe26abe58629623a8b38382b33947c5a5c702ce586e6e6ea51430380bf.json +++ b/backend/.sqlx/query-06072cfe26abe58629623a8b38382b33947c5a5c702ce586e6e6ea51430380bf.json @@ -17,8 +17,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json b/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json index d29a18c691..e7ed0aee65 100644 --- a/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json +++ b/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json @@ -46,11 +46,11 @@ ] }, "nullable": [ - true, - true, - true, - true, - true, + false, + false, + false, + false, + false, true, true ] diff --git a/backend/.sqlx/query-089d7bc7acdbb97cf477159e111bc7e9ee85289ff5c52af43166928337c257e7.json b/backend/.sqlx/query-089d7bc7acdbb97cf477159e111bc7e9ee85289ff5c52af43166928337c257e7.json index b925141065..a032a87239 100644 --- a/backend/.sqlx/query-089d7bc7acdbb97cf477159e111bc7e9ee85289ff5c52af43166928337c257e7.json +++ b/backend/.sqlx/query-089d7bc7acdbb97cf477159e111bc7e9ee85289ff5c52af43166928337c257e7.json @@ -30,8 +30,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-0f7e01b613a94b29784aae6d7b17b23d6dcf2e5364852a5e85b3c41c417bace2.json b/backend/.sqlx/query-0f7e01b613a94b29784aae6d7b17b23d6dcf2e5364852a5e85b3c41c417bace2.json index de4dd4a7e3..1648882c1a 100644 --- a/backend/.sqlx/query-0f7e01b613a94b29784aae6d7b17b23d6dcf2e5364852a5e85b3c41c417bace2.json +++ b/backend/.sqlx/query-0f7e01b613a94b29784aae6d7b17b23d6dcf2e5364852a5e85b3c41c417bace2.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-14276a040cb4db88d71fccdc3579e8c0bb132b70668301b535872d1632753e30.json b/backend/.sqlx/query-14276a040cb4db88d71fccdc3579e8c0bb132b70668301b535872d1632753e30.json index bed99ef1b7..d4f7afa966 100644 --- a/backend/.sqlx/query-14276a040cb4db88d71fccdc3579e8c0bb132b70668301b535872d1632753e30.json +++ b/backend/.sqlx/query-14276a040cb4db88d71fccdc3579e8c0bb132b70668301b535872d1632753e30.json @@ -122,8 +122,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-16d438374b03a9c515f4c2d638366f38ffe2f3a0958adea53e67757c6ac463ec.json b/backend/.sqlx/query-16d438374b03a9c515f4c2d638366f38ffe2f3a0958adea53e67757c6ac463ec.json index 1fdc78c472..c091114374 100644 --- a/backend/.sqlx/query-16d438374b03a9c515f4c2d638366f38ffe2f3a0958adea53e67757c6ac463ec.json +++ b/backend/.sqlx/query-16d438374b03a9c515f4c2d638366f38ffe2f3a0958adea53e67757c6ac463ec.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-19b59c478744d029c6006b01f04243ad2e0aef485a780daea5d76b0be2bb2ea2.json b/backend/.sqlx/query-19b59c478744d029c6006b01f04243ad2e0aef485a780daea5d76b0be2bb2ea2.json index 30e071f8ce..43c68f8c5a 100644 --- a/backend/.sqlx/query-19b59c478744d029c6006b01f04243ad2e0aef485a780daea5d76b0be2bb2ea2.json +++ b/backend/.sqlx/query-19b59c478744d029c6006b01f04243ad2e0aef485a780daea5d76b0be2bb2ea2.json @@ -40,8 +40,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-a7468e9054beed88636786c5495ac3b9d9a6086ae6212ad9237b39a0346d7d26.json b/backend/.sqlx/query-209dc4c1b91eeab1c12ffcd9f9e16f315c689ca772c736b333dcdf07c8086087.json similarity index 53% rename from backend/.sqlx/query-a7468e9054beed88636786c5495ac3b9d9a6086ae6212ad9237b39a0346d7d26.json rename to backend/.sqlx/query-209dc4c1b91eeab1c12ffcd9f9e16f315c689ca772c736b333dcdf07c8086087.json index cead5d7019..e4fbe7efe0 100644 --- a/backend/.sqlx/query-a7468e9054beed88636786c5495ac3b9d9a6086ae6212ad9237b39a0346d7d26.json +++ b/backend/.sqlx/query-209dc4c1b91eeab1c12ffcd9f9e16f315c689ca772c736b333dcdf07c8086087.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT\n flow_version.id AS version,\n flow_version.value->>'early_return' as early_return, \n flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor, \n (flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled, \n flow.tag, \n flow.dedicated_worker, \n flow.on_behalf_of_email, \n flow.edited_by\n FROM \n flow_version\n INNER JOIN flow\n ON flow.path = flow_version.path AND\n flow.workspace_id = flow_version.workspace_id\n WHERE \n flow_version.workspace_id = $1 AND\n flow_version.path = $2 AND\n flow_version.id = $3\n ", + "query": "\n SELECT\n flow_version.id AS version,\n flow_version.value->>'early_return' as early_return,\n flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor,\n (flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled,\n flow.tag,\n flow.dedicated_worker,\n flow.on_behalf_of_email,\n flow.edited_by\n FROM\n flow_version\n INNER JOIN flow\n ON flow.path = flow_version.path AND\n flow.workspace_id = flow_version.workspace_id\n WHERE\n flow_version.workspace_id = $1 AND\n flow_version.path = $2 AND\n flow_version.id = $3\n ", "describe": { "columns": [ { @@ -62,5 +62,5 @@ false ] }, - "hash": "a7468e9054beed88636786c5495ac3b9d9a6086ae6212ad9237b39a0346d7d26" + "hash": "209dc4c1b91eeab1c12ffcd9f9e16f315c689ca772c736b333dcdf07c8086087" } diff --git a/backend/.sqlx/query-212553c83e4dcdc6d045eb2fe2dadbb2860ce52d37a56b2861de1215260ecff8.json b/backend/.sqlx/query-212553c83e4dcdc6d045eb2fe2dadbb2860ce52d37a56b2861de1215260ecff8.json index b0354db034..7d8b2eba1e 100644 --- a/backend/.sqlx/query-212553c83e4dcdc6d045eb2fe2dadbb2860ce52d37a56b2861de1215260ecff8.json +++ b/backend/.sqlx/query-212553c83e4dcdc6d045eb2fe2dadbb2860ce52d37a56b2861de1215260ecff8.json @@ -34,8 +34,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } @@ -68,8 +67,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-22e0e8a1aa48f8b21763452bd36fbe7db4887c4ac5295052c796bd78a7edc50b.json b/backend/.sqlx/query-22e0e8a1aa48f8b21763452bd36fbe7db4887c4ac5295052c796bd78a7edc50b.json index 235f255dd1..66c9e6e96a 100644 --- a/backend/.sqlx/query-22e0e8a1aa48f8b21763452bd36fbe7db4887c4ac5295052c796bd78a7edc50b.json +++ b/backend/.sqlx/query-22e0e8a1aa48f8b21763452bd36fbe7db4887c4ac5295052c796bd78a7edc50b.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-23419adcd74c326d716527293eff518b42f4cdb33e034441015494bd26c172d2.json b/backend/.sqlx/query-23419adcd74c326d716527293eff518b42f4cdb33e034441015494bd26c172d2.json index 7bd7367d8d..242b358ff5 100644 --- a/backend/.sqlx/query-23419adcd74c326d716527293eff518b42f4cdb33e034441015494bd26c172d2.json +++ b/backend/.sqlx/query-23419adcd74c326d716527293eff518b42f4cdb33e034441015494bd26c172d2.json @@ -40,8 +40,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-256118bb9c87675d6a8db2d13c4f446d88c7bc671d5e90307febfd60bb235031.json b/backend/.sqlx/query-256118bb9c87675d6a8db2d13c4f446d88c7bc671d5e90307febfd60bb235031.json index c51175e921..4e48ae281c 100644 --- a/backend/.sqlx/query-256118bb9c87675d6a8db2d13c4f446d88c7bc671d5e90307febfd60bb235031.json +++ b/backend/.sqlx/query-256118bb9c87675d6a8db2d13c4f446d88c7bc671d5e90307febfd60bb235031.json @@ -16,8 +16,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-2602402bedcdbc45cdc0d64a76ce075c6ae51404037b0a7d4d33faf6a7d6a6d8.json b/backend/.sqlx/query-2602402bedcdbc45cdc0d64a76ce075c6ae51404037b0a7d4d33faf6a7d6a6d8.json index 5a4cd89224..128a435c8b 100644 --- a/backend/.sqlx/query-2602402bedcdbc45cdc0d64a76ce075c6ae51404037b0a7d4d33faf6a7d6a6d8.json +++ b/backend/.sqlx/query-2602402bedcdbc45cdc0d64a76ce075c6ae51404037b0a7d4d33faf6a7d6a6d8.json @@ -11,8 +11,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-26141a1c3c48184285d94a2e2aca94ff708ce8e0b0b4fb59e5de7c9698c75c25.json b/backend/.sqlx/query-26141a1c3c48184285d94a2e2aca94ff708ce8e0b0b4fb59e5de7c9698c75c25.json index 08c6c1cbc2..9abe8af9af 100644 --- a/backend/.sqlx/query-26141a1c3c48184285d94a2e2aca94ff708ce8e0b0b4fb59e5de7c9698c75c25.json +++ b/backend/.sqlx/query-26141a1c3c48184285d94a2e2aca94ff708ce8e0b0b4fb59e5de7c9698c75c25.json @@ -11,8 +11,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-27ada97cb533c8595f1d73987c7823d8e54c96889e06895c57cafae9ca27bf8b.json b/backend/.sqlx/query-27ada97cb533c8595f1d73987c7823d8e54c96889e06895c57cafae9ca27bf8b.json index 1199e9441f..5d7c0e2abc 100644 --- a/backend/.sqlx/query-27ada97cb533c8595f1d73987c7823d8e54c96889e06895c57cafae9ca27bf8b.json +++ b/backend/.sqlx/query-27ada97cb533c8595f1d73987c7823d8e54c96889e06895c57cafae9ca27bf8b.json @@ -15,8 +15,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-2fd22c4ffa2d222bb116260994a748e0639c2f73cbc1d8be66420c70b14c96e1.json b/backend/.sqlx/query-2fd22c4ffa2d222bb116260994a748e0639c2f73cbc1d8be66420c70b14c96e1.json index afd0f503bf..8f6c8edf25 100644 --- a/backend/.sqlx/query-2fd22c4ffa2d222bb116260994a748e0639c2f73cbc1d8be66420c70b14c96e1.json +++ b/backend/.sqlx/query-2fd22c4ffa2d222bb116260994a748e0639c2f73cbc1d8be66420c70b14c96e1.json @@ -12,8 +12,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-3b3f60623126626b52ca0a4a188655ddf728cd3f21ee308db7393694ccc5c7b3.json b/backend/.sqlx/query-3b3f60623126626b52ca0a4a188655ddf728cd3f21ee308db7393694ccc5c7b3.json index 0df550b57c..27b2503991 100644 --- a/backend/.sqlx/query-3b3f60623126626b52ca0a4a188655ddf728cd3f21ee308db7393694ccc5c7b3.json +++ b/backend/.sqlx/query-3b3f60623126626b52ca0a4a188655ddf728cd3f21ee308db7393694ccc5c7b3.json @@ -12,8 +12,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-3e5bdc2e071fc2f1e3c7971736272f20bb5a0aa921a614bd02898d3f162660c2.json b/backend/.sqlx/query-3e5bdc2e071fc2f1e3c7971736272f20bb5a0aa921a614bd02898d3f162660c2.json index fdcb10be1c..e6dc6ff08b 100644 --- a/backend/.sqlx/query-3e5bdc2e071fc2f1e3c7971736272f20bb5a0aa921a614bd02898d3f162660c2.json +++ b/backend/.sqlx/query-3e5bdc2e071fc2f1e3c7971736272f20bb5a0aa921a614bd02898d3f162660c2.json @@ -16,8 +16,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } @@ -52,8 +51,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-42b4b73e9d60348e2d90fcade9dcad6d8995242dc20a4e14c1a8fae4fc6a9fd2.json b/backend/.sqlx/query-42b4b73e9d60348e2d90fcade9dcad6d8995242dc20a4e14c1a8fae4fc6a9fd2.json index 537ecbeacf..e6d71f386f 100644 --- a/backend/.sqlx/query-42b4b73e9d60348e2d90fcade9dcad6d8995242dc20a4e14c1a8fae4fc6a9fd2.json +++ b/backend/.sqlx/query-42b4b73e9d60348e2d90fcade9dcad6d8995242dc20a4e14c1a8fae4fc6a9fd2.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-45024b932383199974616bba1fc2f7175cc6f2e02d9c565bb5159cae3e0b6835.json b/backend/.sqlx/query-45024b932383199974616bba1fc2f7175cc6f2e02d9c565bb5159cae3e0b6835.json index 30c5ff7a49..f215cf6147 100644 --- a/backend/.sqlx/query-45024b932383199974616bba1fc2f7175cc6f2e02d9c565bb5159cae3e0b6835.json +++ b/backend/.sqlx/query-45024b932383199974616bba1fc2f7175cc6f2e02d9c565bb5159cae3e0b6835.json @@ -30,8 +30,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-4d272cf4a77aab7007a5b35589e08532a1020cabaf5e22325a1e05f0491d785c.json b/backend/.sqlx/query-4d272cf4a77aab7007a5b35589e08532a1020cabaf5e22325a1e05f0491d785c.json index cfb975b838..896b21dc7f 100644 --- a/backend/.sqlx/query-4d272cf4a77aab7007a5b35589e08532a1020cabaf5e22325a1e05f0491d785c.json +++ b/backend/.sqlx/query-4d272cf4a77aab7007a5b35589e08532a1020cabaf5e22325a1e05f0491d785c.json @@ -37,8 +37,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-4f547c0fd54f3bc57212ce87810e35adf640d44d607e62a1fb296e38ac3fdd36.json b/backend/.sqlx/query-4f547c0fd54f3bc57212ce87810e35adf640d44d607e62a1fb296e38ac3fdd36.json index 4e8ec56d12..789c0334d7 100644 --- a/backend/.sqlx/query-4f547c0fd54f3bc57212ce87810e35adf640d44d607e62a1fb296e38ac3fdd36.json +++ b/backend/.sqlx/query-4f547c0fd54f3bc57212ce87810e35adf640d44d607e62a1fb296e38ac3fdd36.json @@ -32,8 +32,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } @@ -71,8 +70,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-5368683c19f8d6744d5dbc53e5b2ab0f2348646d79f5306c6868e2c3a8f389ee.json b/backend/.sqlx/query-5368683c19f8d6744d5dbc53e5b2ab0f2348646d79f5306c6868e2c3a8f389ee.json index ebf1df39a3..1058d78d8b 100644 --- a/backend/.sqlx/query-5368683c19f8d6744d5dbc53e5b2ab0f2348646d79f5306c6868e2c3a8f389ee.json +++ b/backend/.sqlx/query-5368683c19f8d6744d5dbc53e5b2ab0f2348646d79f5306c6868e2c3a8f389ee.json @@ -16,8 +16,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-5bf200f2c8db25ddf231b564503c6c70f7f3958564a79bb0c6b3863b1ebb0cbf.json b/backend/.sqlx/query-5bf200f2c8db25ddf231b564503c6c70f7f3958564a79bb0c6b3863b1ebb0cbf.json index 9422e3e8d0..6d3941ad69 100644 --- a/backend/.sqlx/query-5bf200f2c8db25ddf231b564503c6c70f7f3958564a79bb0c6b3863b1ebb0cbf.json +++ b/backend/.sqlx/query-5bf200f2c8db25ddf231b564503c6c70f7f3958564a79bb0c6b3863b1ebb0cbf.json @@ -245,8 +245,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-66a0e51cf149ba532463e29dd361a803e1bced2f8e1a12f8933b7598ee85a147.json b/backend/.sqlx/query-66a0e51cf149ba532463e29dd361a803e1bced2f8e1a12f8933b7598ee85a147.json index 5f9be1cba9..0d12be2448 100644 --- a/backend/.sqlx/query-66a0e51cf149ba532463e29dd361a803e1bced2f8e1a12f8933b7598ee85a147.json +++ b/backend/.sqlx/query-66a0e51cf149ba532463e29dd361a803e1bced2f8e1a12f8933b7598ee85a147.json @@ -35,8 +35,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-7065f23d04e26831664048f2cfc4f412c57af931f80621aee5012e9cb3535626.json b/backend/.sqlx/query-7065f23d04e26831664048f2cfc4f412c57af931f80621aee5012e9cb3535626.json index 48ea42f8a2..02823e8291 100644 --- a/backend/.sqlx/query-7065f23d04e26831664048f2cfc4f412c57af931f80621aee5012e9cb3535626.json +++ b/backend/.sqlx/query-7065f23d04e26831664048f2cfc4f412c57af931f80621aee5012e9cb3535626.json @@ -29,8 +29,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-73fdd01bad58b8be1a52f89faef8d92a983470adcd3cc850734960c905e61e83.json b/backend/.sqlx/query-73fdd01bad58b8be1a52f89faef8d92a983470adcd3cc850734960c905e61e83.json index 8aacd8a805..cba7ffdfef 100644 --- a/backend/.sqlx/query-73fdd01bad58b8be1a52f89faef8d92a983470adcd3cc850734960c905e61e83.json +++ b/backend/.sqlx/query-73fdd01bad58b8be1a52f89faef8d92a983470adcd3cc850734960c905e61e83.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-757ef6215d3d385cb3a69e26ee4ca846dd5e7fe7ceb1aa8b3fcd26a2bd30eb2c.json b/backend/.sqlx/query-757ef6215d3d385cb3a69e26ee4ca846dd5e7fe7ceb1aa8b3fcd26a2bd30eb2c.json index 59b56ceda7..cd795e6fec 100644 --- a/backend/.sqlx/query-757ef6215d3d385cb3a69e26ee4ca846dd5e7fe7ceb1aa8b3fcd26a2bd30eb2c.json +++ b/backend/.sqlx/query-757ef6215d3d385cb3a69e26ee4ca846dd5e7fe7ceb1aa8b3fcd26a2bd30eb2c.json @@ -40,8 +40,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-7fbf72d9059fcd77e4c1112fa4fa22e4276c1da653475628889ce17dc904fbaa.json b/backend/.sqlx/query-7fbf72d9059fcd77e4c1112fa4fa22e4276c1da653475628889ce17dc904fbaa.json index 8c380df861..393a920b7c 100644 --- a/backend/.sqlx/query-7fbf72d9059fcd77e4c1112fa4fa22e4276c1da653475628889ce17dc904fbaa.json +++ b/backend/.sqlx/query-7fbf72d9059fcd77e4c1112fa4fa22e4276c1da653475628889ce17dc904fbaa.json @@ -27,8 +27,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-836bac47d89113d90bd03a471446eb9016207975af1e37042d81df8cb6ae2c53.json b/backend/.sqlx/query-836bac47d89113d90bd03a471446eb9016207975af1e37042d81df8cb6ae2c53.json index 3012e9ef77..437d644eb2 100644 --- a/backend/.sqlx/query-836bac47d89113d90bd03a471446eb9016207975af1e37042d81df8cb6ae2c53.json +++ b/backend/.sqlx/query-836bac47d89113d90bd03a471446eb9016207975af1e37042d81df8cb6ae2c53.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-87564a196a1662f524407d853db506bf08c28efe82b68b3d44bafbd3d0e91c29.json b/backend/.sqlx/query-87564a196a1662f524407d853db506bf08c28efe82b68b3d44bafbd3d0e91c29.json index 5a61f82be3..9350442134 100644 --- a/backend/.sqlx/query-87564a196a1662f524407d853db506bf08c28efe82b68b3d44bafbd3d0e91c29.json +++ b/backend/.sqlx/query-87564a196a1662f524407d853db506bf08c28efe82b68b3d44bafbd3d0e91c29.json @@ -35,8 +35,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-9052b7cd438ff029a37bd489190d98d365acec09f2f102b7de71dcc9d356900e.json b/backend/.sqlx/query-9052b7cd438ff029a37bd489190d98d365acec09f2f102b7de71dcc9d356900e.json index 62f8d3100f..02cad7d11d 100644 --- a/backend/.sqlx/query-9052b7cd438ff029a37bd489190d98d365acec09f2f102b7de71dcc9d356900e.json +++ b/backend/.sqlx/query-9052b7cd438ff029a37bd489190d98d365acec09f2f102b7de71dcc9d356900e.json @@ -17,8 +17,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-940b6d78bab940a37a42492f030d2393e297043e4e58555d872b5c4dd89c196a.json b/backend/.sqlx/query-940b6d78bab940a37a42492f030d2393e297043e4e58555d872b5c4dd89c196a.json index 09aadbda1e..47554ad43e 100644 --- a/backend/.sqlx/query-940b6d78bab940a37a42492f030d2393e297043e4e58555d872b5c4dd89c196a.json +++ b/backend/.sqlx/query-940b6d78bab940a37a42492f030d2393e297043e4e58555d872b5c4dd89c196a.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-95c57fb921a2e3725b92cbafac6e3dc360b88429f03dd1e2b1b55cfabe208cb7.json b/backend/.sqlx/query-95c57fb921a2e3725b92cbafac6e3dc360b88429f03dd1e2b1b55cfabe208cb7.json index a189bc4da2..4b989c893b 100644 --- a/backend/.sqlx/query-95c57fb921a2e3725b92cbafac6e3dc360b88429f03dd1e2b1b55cfabe208cb7.json +++ b/backend/.sqlx/query-95c57fb921a2e3725b92cbafac6e3dc360b88429f03dd1e2b1b55cfabe208cb7.json @@ -17,8 +17,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-9c50e3a136a8ee3ec56e083f26d3a960b89e02ec40b292f3b5198baf2a1d3dbf.json b/backend/.sqlx/query-9c50e3a136a8ee3ec56e083f26d3a960b89e02ec40b292f3b5198baf2a1d3dbf.json index 66a487695f..9759bad4d4 100644 --- a/backend/.sqlx/query-9c50e3a136a8ee3ec56e083f26d3a960b89e02ec40b292f3b5198baf2a1d3dbf.json +++ b/backend/.sqlx/query-9c50e3a136a8ee3ec56e083f26d3a960b89e02ec40b292f3b5198baf2a1d3dbf.json @@ -32,8 +32,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-9ecb404e46a4eac55f977f05a3afbafe5dc3cdecc17a3d5a7476b160c1b6e7e1.json b/backend/.sqlx/query-9ecb404e46a4eac55f977f05a3afbafe5dc3cdecc17a3d5a7476b160c1b6e7e1.json index b6fe1b7fa0..79c8f0b45d 100644 --- a/backend/.sqlx/query-9ecb404e46a4eac55f977f05a3afbafe5dc3cdecc17a3d5a7476b160c1b6e7e1.json +++ b/backend/.sqlx/query-9ecb404e46a4eac55f977f05a3afbafe5dc3cdecc17a3d5a7476b160c1b6e7e1.json @@ -30,8 +30,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-a1745a4f525b251d2f5a602ab2b2ede46b4471e21b11f607573a844013911abe.json b/backend/.sqlx/query-a1745a4f525b251d2f5a602ab2b2ede46b4471e21b11f607573a844013911abe.json index 33636da608..a25f845f91 100644 --- a/backend/.sqlx/query-a1745a4f525b251d2f5a602ab2b2ede46b4471e21b11f607573a844013911abe.json +++ b/backend/.sqlx/query-a1745a4f525b251d2f5a602ab2b2ede46b4471e21b11f607573a844013911abe.json @@ -155,8 +155,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-a4b6371d33206010b2f3ffd2b09e33244fe8ab9a803248fc23f334034d24aad4.json b/backend/.sqlx/query-a4b6371d33206010b2f3ffd2b09e33244fe8ab9a803248fc23f334034d24aad4.json index 03ee58ca3f..da6d213748 100644 --- a/backend/.sqlx/query-a4b6371d33206010b2f3ffd2b09e33244fe8ab9a803248fc23f334034d24aad4.json +++ b/backend/.sqlx/query-a4b6371d33206010b2f3ffd2b09e33244fe8ab9a803248fc23f334034d24aad4.json @@ -185,8 +185,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-b3771b690c5966272b1f42c9965bb6a8f961c119516e4c33dc928cd3b4f4edbc.json b/backend/.sqlx/query-b3771b690c5966272b1f42c9965bb6a8f961c119516e4c33dc928cd3b4f4edbc.json index 4fa7da00e0..d08af6ffdd 100644 --- a/backend/.sqlx/query-b3771b690c5966272b1f42c9965bb6a8f961c119516e4c33dc928cd3b4f4edbc.json +++ b/backend/.sqlx/query-b3771b690c5966272b1f42c9965bb6a8f961c119516e4c33dc928cd3b4f4edbc.json @@ -160,8 +160,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-b3f0595cacba194e08b9a3e244d9e637e9e156cd85b69126c87dfff89a47711d.json b/backend/.sqlx/query-b3f0595cacba194e08b9a3e244d9e637e9e156cd85b69126c87dfff89a47711d.json index 186822a000..54a4e3cd93 100644 --- a/backend/.sqlx/query-b3f0595cacba194e08b9a3e244d9e637e9e156cd85b69126c87dfff89a47711d.json +++ b/backend/.sqlx/query-b3f0595cacba194e08b9a3e244d9e637e9e156cd85b69126c87dfff89a47711d.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-b615d73ddb43e9d655b86a0cf98f892bf40e629ee11ee4845199481755f2789d.json b/backend/.sqlx/query-b615d73ddb43e9d655b86a0cf98f892bf40e629ee11ee4845199481755f2789d.json index e6ef13f492..96f59aab3f 100644 --- a/backend/.sqlx/query-b615d73ddb43e9d655b86a0cf98f892bf40e629ee11ee4845199481755f2789d.json +++ b/backend/.sqlx/query-b615d73ddb43e9d655b86a0cf98f892bf40e629ee11ee4845199481755f2789d.json @@ -21,8 +21,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } @@ -72,8 +71,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-be6d2c92a62b7b284651c45af809746147aa9b8d0a81642a7b7cb4738a0cad66.json b/backend/.sqlx/query-be6d2c92a62b7b284651c45af809746147aa9b8d0a81642a7b7cb4738a0cad66.json index 4118760af2..0688afebd3 100644 --- a/backend/.sqlx/query-be6d2c92a62b7b284651c45af809746147aa9b8d0a81642a7b7cb4738a0cad66.json +++ b/backend/.sqlx/query-be6d2c92a62b7b284651c45af809746147aa9b8d0a81642a7b7cb4738a0cad66.json @@ -105,8 +105,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-c3b1152b554812d65eb27f95b1fd434f860922fbc021185beffb9827647feb8e.json b/backend/.sqlx/query-c3b1152b554812d65eb27f95b1fd434f860922fbc021185beffb9827647feb8e.json index 07243717a0..0098e51ab2 100644 --- a/backend/.sqlx/query-c3b1152b554812d65eb27f95b1fd434f860922fbc021185beffb9827647feb8e.json +++ b/backend/.sqlx/query-c3b1152b554812d65eb27f95b1fd434f860922fbc021185beffb9827647feb8e.json @@ -31,8 +31,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-cc5919b087e3319045fc2242f22384514987b6a4eb807ee375fb8b5cd4e92307.json b/backend/.sqlx/query-cc5919b087e3319045fc2242f22384514987b6a4eb807ee375fb8b5cd4e92307.json index 8a0dacd864..0e338633fe 100644 --- a/backend/.sqlx/query-cc5919b087e3319045fc2242f22384514987b6a4eb807ee375fb8b5cd4e92307.json +++ b/backend/.sqlx/query-cc5919b087e3319045fc2242f22384514987b6a4eb807ee375fb8b5cd4e92307.json @@ -11,8 +11,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-ccef7a1bde5cac6c362c5fedb6c13f1f882b695f896f94e5cf91d205633355a1.json b/backend/.sqlx/query-ccef7a1bde5cac6c362c5fedb6c13f1f882b695f896f94e5cf91d205633355a1.json index 2b39145512..5091b2fc69 100644 --- a/backend/.sqlx/query-ccef7a1bde5cac6c362c5fedb6c13f1f882b695f896f94e5cf91d205633355a1.json +++ b/backend/.sqlx/query-ccef7a1bde5cac6c362c5fedb6c13f1f882b695f896f94e5cf91d205633355a1.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-d0b493da8ff5b3b78a9a0c6972cd2ca5449dad0fe361bee1a9834989bfeabbae.json b/backend/.sqlx/query-d0b493da8ff5b3b78a9a0c6972cd2ca5449dad0fe361bee1a9834989bfeabbae.json index b71bf15f85..b0cef08da3 100644 --- a/backend/.sqlx/query-d0b493da8ff5b3b78a9a0c6972cd2ca5449dad0fe361bee1a9834989bfeabbae.json +++ b/backend/.sqlx/query-d0b493da8ff5b3b78a9a0c6972cd2ca5449dad0fe361bee1a9834989bfeabbae.json @@ -12,8 +12,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-d41ea93fd58381b89e151c965eae1ea2fe96a1b94f5a92953fb1c1642d15c016.json b/backend/.sqlx/query-d41ea93fd58381b89e151c965eae1ea2fe96a1b94f5a92953fb1c1642d15c016.json index 9d2f8d9a9f..11271e94d8 100644 --- a/backend/.sqlx/query-d41ea93fd58381b89e151c965eae1ea2fe96a1b94f5a92953fb1c1642d15c016.json +++ b/backend/.sqlx/query-d41ea93fd58381b89e151c965eae1ea2fe96a1b94f5a92953fb1c1642d15c016.json @@ -105,8 +105,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-d495c94b580fd34d5ae90615ef21a8a9cc35f362197c0766a5787436af141106.json b/backend/.sqlx/query-d495c94b580fd34d5ae90615ef21a8a9cc35f362197c0766a5787436af141106.json index 67a38d1904..9cef6500d3 100644 --- a/backend/.sqlx/query-d495c94b580fd34d5ae90615ef21a8a9cc35f362197c0766a5787436af141106.json +++ b/backend/.sqlx/query-d495c94b580fd34d5ae90615ef21a8a9cc35f362197c0766a5787436af141106.json @@ -25,8 +25,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-e4d71278fb80126a7a9da73f1889352d4d1e3cb3a8a08f1c9c03055a1cab1235.json b/backend/.sqlx/query-e4d71278fb80126a7a9da73f1889352d4d1e3cb3a8a08f1c9c03055a1cab1235.json index 7309b03a02..5b07bcd9c9 100644 --- a/backend/.sqlx/query-e4d71278fb80126a7a9da73f1889352d4d1e3cb3a8a08f1c9c03055a1cab1235.json +++ b/backend/.sqlx/query-e4d71278fb80126a7a9da73f1889352d4d1e3cb3a8a08f1c9c03055a1cab1235.json @@ -185,8 +185,7 @@ "sqs", "gcp", "mqtt", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-e80177f3ffd4c1f52cdb4757483f03f72ef81db302d727e18e63a307ac902022.json b/backend/.sqlx/query-e80177f3ffd4c1f52cdb4757483f03f72ef81db302d727e18e63a307ac902022.json index dd65a586d1..1e38e57fb5 100644 --- a/backend/.sqlx/query-e80177f3ffd4c1f52cdb4757483f03f72ef81db302d727e18e63a307ac902022.json +++ b/backend/.sqlx/query-e80177f3ffd4c1f52cdb4757483f03f72ef81db302d727e18e63a307ac902022.json @@ -31,8 +31,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-eac595e19e5c8e70f1514ef29dec35c7342ac9a814c73f6290e1d6ebd3a55423.json b/backend/.sqlx/query-eac595e19e5c8e70f1514ef29dec35c7342ac9a814c73f6290e1d6ebd3a55423.json index 7c45c44b13..57607ae052 100644 --- a/backend/.sqlx/query-eac595e19e5c8e70f1514ef29dec35c7342ac9a814c73f6290e1d6ebd3a55423.json +++ b/backend/.sqlx/query-eac595e19e5c8e70f1514ef29dec35c7342ac9a814c73f6290e1d6ebd3a55423.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-ecab1af12a7afa685c056b9d0e526275203fc8ecddf83ca6d05c9fb77e46e7ee.json b/backend/.sqlx/query-ecab1af12a7afa685c056b9d0e526275203fc8ecddf83ca6d05c9fb77e46e7ee.json index 16ccdd10e9..f80b75b5ec 100644 --- a/backend/.sqlx/query-ecab1af12a7afa685c056b9d0e526275203fc8ecddf83ca6d05c9fb77e46e7ee.json +++ b/backend/.sqlx/query-ecab1af12a7afa685c056b9d0e526275203fc8ecddf83ca6d05c9fb77e46e7ee.json @@ -21,8 +21,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } @@ -72,8 +71,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-ed8facbf29ebb670d05fe8aa34b50d6a6935420fbedc83aa3ad1e9be7465c8dd.json b/backend/.sqlx/query-ed8facbf29ebb670d05fe8aa34b50d6a6935420fbedc83aa3ad1e9be7465c8dd.json index 1876db927e..43fd90abe5 100644 --- a/backend/.sqlx/query-ed8facbf29ebb670d05fe8aa34b50d6a6935420fbedc83aa3ad1e9be7465c8dd.json +++ b/backend/.sqlx/query-ed8facbf29ebb670d05fe8aa34b50d6a6935420fbedc83aa3ad1e9be7465c8dd.json @@ -24,8 +24,7 @@ "mqtt", "gcp", "default_email", - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/.sqlx/query-ee537def1ead8bee48bb9f5c1f57d42e7add6011c34d91761ba23e2c74c4032c.json b/backend/.sqlx/query-ee537def1ead8bee48bb9f5c1f57d42e7add6011c34d91761ba23e2c74c4032c.json index 536175599d..4ce08c9e1c 100644 --- a/backend/.sqlx/query-ee537def1ead8bee48bb9f5c1f57d42e7add6011c34d91761ba23e2c74c4032c.json +++ b/backend/.sqlx/query-ee537def1ead8bee48bb9f5c1f57d42e7add6011c34d91761ba23e2c74c4032c.json @@ -21,8 +21,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } @@ -72,8 +71,7 @@ "name": "native_trigger_service", "kind": { "Enum": [ - "nextcloud", - "google" + "nextcloud" ] } } diff --git a/backend/src/db_connect.rs b/backend/src/db_connect.rs new file mode 100644 index 0000000000..05bf7b6522 --- /dev/null +++ b/backend/src/db_connect.rs @@ -0,0 +1,159 @@ +use windmill_common::{ + error::{self, Error}, + get_database_url, DatabaseUrl, +}; + +pub const DEFAULT_MAX_CONNECTIONS_SERVER: u32 = 50; +pub const DEFAULT_MAX_CONNECTIONS_WORKER: u32 = 5; +pub const DEFAULT_MAX_CONNECTIONS_INDEXER: u32 = 5; + +pub async fn initial_connection() -> Result, error::Error> { + let connect_options = get_database_url().await?.connect_options().await?; + sqlx::postgres::PgPoolOptions::new() + .max_connections(2) + .connect_with(connect_options) + .await + .map_err(|err| Error::ConnectingToDatabase(err.to_string())) +} + +pub async fn connect_db( + server_mode: bool, + indexer_mode: bool, + worker_mode: bool, + #[cfg(feature = "private")] mut killpill_rx: tokio::sync::broadcast::Receiver<()>, +) -> anyhow::Result> { + use anyhow::Context; + + let database_url = get_database_url().await?; + + let max_connections = match std::env::var("DATABASE_CONNECTIONS") { + Ok(n) => n.parse::().context("invalid DATABASE_CONNECTIONS")?, + Err(_) => { + if server_mode { + DEFAULT_MAX_CONNECTIONS_SERVER + } else if indexer_mode { + DEFAULT_MAX_CONNECTIONS_INDEXER + } else { + DEFAULT_MAX_CONNECTIONS_WORKER + + std::env::var("NUM_WORKERS") + .ok() + .map(|x| x.parse().ok()) + .flatten() + .unwrap_or(1) + - 1 + } + } + }; + + let pool = connect(database_url.clone(), max_connections, worker_mode).await?; + #[cfg(all(feature = "enterprise", feature = "private"))] + let pool2 = pool.clone(); + #[cfg(all(feature = "enterprise", feature = "private"))] + if let DatabaseUrl::IamRds(database_url) = database_url { + tokio::spawn(async move { + loop { + tokio::select! { + _ = killpill_rx.recv() => { + break; + } + _ = tokio::time::sleep(std::time::Duration::from_secs(10)) => { + let needs_refresh = { + let read_guard = database_url.read().await; + read_guard.needs_refresh() + }; + if needs_refresh { + let new_url = tokio::time::timeout(std::time::Duration::from_secs(10), get_database_url()).await; + match new_url { + Ok(Ok(new_url)) => { + match new_url.connect_options().await { + Ok(connect_options) => { + pool2.set_connect_options(connect_options); + tracing::info!("Refreshed IAM RDS URL successfully"); + } + Err(e) => { + tracing::error!("Error getting IAM RDS connect options, retrying in 10s: {}", e); + continue; + } + } + } + Ok(Err(e)) => { + tracing::error!("Error refreshing IAM RDS URL, trying again in 10s: {}", e); + continue; + } + Err(e) => { + tracing::error!("Timeout after 10s refreshing IAM RDS URL, trying again in 10 seconds: {}", e); + continue; + } + } + } + } + } + } + }); + } + + Ok(pool) +} + +pub async fn connect( + database_url: DatabaseUrl, + max_connections: u32, + worker_mode: bool, +) -> Result, error::Error> { + use sqlx::Executor; + use std::time::Duration; + let mut pool_options = sqlx::postgres::PgPoolOptions::new() + .min_connections((max_connections / 5).clamp(1, max_connections)) + .max_connections(max_connections) + .max_lifetime(Duration::from_secs(30 * 60)); // 30 mins + if worker_mode { + pool_options = pool_options.idle_timeout(Duration::from_secs(60)); + } + pool_options + .after_connect(move |conn, _| { + if worker_mode { + Box::pin(async move { + if let Err(e) = conn + .execute( + r#" + SET enable_seqscan = OFF; + SET statement_timeout = '5min'; + SET idle_in_transaction_session_timeout = '10min'; + SET tcp_keepalives_idle = 300; + SET tcp_keepalives_interval = 60; + SET tcp_keepalives_count = 10;"#, + ) + .await + { + tracing::error!("Error setting postgres settings: {}", e); + } + Ok(()) + }) + } else { + Box::pin(async move { + if let Err(e) = conn + .execute( + r#" + SET statement_timeout = '5min'; + SET idle_in_transaction_session_timeout = '10min'; + SET tcp_keepalives_idle = 300; + SET tcp_keepalives_interval = 60; + SET tcp_keepalives_count = 10;"#, + ) + .await + { + tracing::error!("Error setting postgres settings: {}", e); + } + Ok(()) + }) + } + }) + .connect_with( + database_url + .connect_options() + .await? + .statement_cache_capacity(400), + ) + .await + .map_err(|err| Error::ConnectingToDatabase(err.to_string())) +} diff --git a/backend/src/main.rs b/backend/src/main.rs index be424a6712..08f6d081a5 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -117,6 +117,7 @@ const BIND_ADDR_ENV: &str = "SERVER_BIND_ADDR"; #[cfg(target_os = "linux")] mod cgroups; +mod db_connect; #[cfg(feature = "private")] pub mod ee; mod ee_oss; @@ -665,7 +666,7 @@ async fn windmill_main() -> anyhow::Result<()> { } else { println!("Connecting to database..."); - let db = windmill_common::initial_connection().await?; + let db = crate::db_connect::initial_connection().await?; let num_version = sqlx::query_scalar!("SELECT version()").fetch_one(&db).await; @@ -770,8 +771,12 @@ async fn windmill_main() -> anyhow::Result<()> { let conn = if mode == Mode::Agent { conn } else { - // This time we use a pool of connections - let db = windmill_common::connect_db( + // Drop the initial connection pool before creating the main one. + // With low PostgreSQL max_connections, both pools existing simultaneously + // can exhaust all available connection slots, causing connect_db to hang. + drop(conn); + + let db = crate::db_connect::connect_db( server_mode, indexer_mode, worker_mode, diff --git a/backend/update_sqlx.sh b/backend/update_sqlx.sh index 669aa90dc9..f87fd72d5c 100755 --- a/backend/update_sqlx.sh +++ b/backend/update_sqlx.sh @@ -10,22 +10,11 @@ if [[ "$(uname)" == "Darwin" ]]; then # Uncomment the git-based samael dependency sed -i '' 's/^# \(samael = { git="https:\/\/github.com\/njaremko\/samael", rev="464d015e3ae393e4b5dd00b4d6baa1b617de0dd6", features = \["xmlsec"\] }\)/\1/' Cargo.toml - # Run cargo sqlx prepare with deno_core_mac + # Run cargo sqlx prepare with deno_core_mac echo "Running cargo sqlx prepare with deno_core_mac..." cargo sqlx prepare --workspace -- --all-targets --features all_sqlx_features,private,deno_core_mac else # Run cargo sqlx prepare echo "Running cargo sqlx prepare..." - cargo sqlx prepare --workspace -- --all-targets --features all_sqlx_features,private -fi - - - -# Undo the samael changes on macOS -if [[ "$(uname)" == "Darwin" ]]; then - echo "Reverting samael changes..." - # Uncomment the version-based samael dependency - sed -i '' 's/^#samael = { version="0.0.14", features = \["xmlsec"\] }/samael = { version="0.0.14", features = ["xmlsec"] }/' Cargo.toml - # Comment out the git-based samael dependency - sed -i '' 's/^\(samael = { git="https:\/\/github.com\/njaremko\/samael", rev="464d015e3ae393e4b5dd00b4d6baa1b617de0dd6", features = \["xmlsec"\] }\)/# \1/' Cargo.toml + cargo sqlx prepare --workspace -- --all-targets --features all_sqlx_features,ee fi diff --git a/backend/v8.snap b/backend/v8.snap deleted file mode 100644 index 4213db10e7..0000000000 Binary files a/backend/v8.snap and /dev/null differ diff --git a/backend/windmill-common/src/bench.rs b/backend/windmill-common/src/bench.rs index cd475edaac..00cb565029 100644 --- a/backend/windmill-common/src/bench.rs +++ b/backend/windmill-common/src/bench.rs @@ -3,7 +3,20 @@ use crate::{ DB, }; use serde::Serialize; +use std::collections::HashSet; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::Arc; use tokio::time::Instant; +use uuid::Uuid; + +static BENCHMARK_INITIALIZED: AtomicBool = AtomicBool::new(false); + +static SHARED_BENCH_ITERS: std::sync::LazyLock> = + std::sync::LazyLock::new(|| Arc::new(AtomicU64::new(0))); + +pub fn shared_bench_iters() -> Arc { + SHARED_BENCH_ITERS.clone() +} #[derive(Serialize)] pub struct PoolStats { @@ -36,6 +49,10 @@ pub struct BenchmarkInfo { pub start: Instant, #[serde(skip)] pub iters: u64, + #[serde(skip)] + seen_top_level: HashSet, + #[serde(skip)] + pub shared_iters: Arc, timings: Vec, pub iter_durations: Vec, pub total_duration: Option, @@ -43,9 +60,11 @@ pub struct BenchmarkInfo { } impl BenchmarkInfo { - pub fn new() -> Self { + pub fn new(shared_iters: Arc) -> Self { BenchmarkInfo { iters: 0, + seen_top_level: HashSet::new(), + shared_iters, timings: vec![], start: Instant::now(), iter_durations: vec![], @@ -64,13 +83,23 @@ impl BenchmarkInfo { } } - pub fn add_iter(&mut self, bench: BenchmarkIter, inc_iters: bool) { - if inc_iters { + pub fn count_top_level(&mut self, job_id: Uuid) -> bool { + if self.seen_top_level.insert(job_id) { + self.iters += 1; + return true; + } + false + } + + pub fn add_iter(&mut self, bench: BenchmarkIter, job_id: Uuid, is_top_level: bool) -> bool { + let newly_counted = is_top_level && self.seen_top_level.insert(job_id); + if newly_counted { self.iters += 1; } let elapsed_total = bench.start.elapsed().as_nanos() as u64; self.timings.push(bench); self.iter_durations.push(elapsed_total); + newly_counted } pub fn write_to_file(&mut self, path: &str) -> anyhow::Result<()> { @@ -118,12 +147,120 @@ impl BenchmarkIter { } } +pub async fn benchmark_verify(benchmark_jobs: i32, db: &DB) { + let benchmark_kind = std::env::var("BENCHMARK_KIND").unwrap_or("noop".to_string()); + + if benchmark_jobs <= 0 || benchmark_kind == "none" { + return; + } + + // For flows, child jobs are created dynamically so only check top-level (parent_job IS NULL). + // "parallelflow" inserts only 1 top-level flow regardless of benchmark_jobs. + let expected_top_level = match benchmark_kind.as_str() { + "parallelflow" => 1i64, + _ => benchmark_jobs as i64, + }; + + let row = sqlx::query!( + "SELECT + COUNT(*) FILTER (WHERE status = 'success') AS succeeded, + COUNT(*) FILTER (WHERE status = 'failure') AS failed, + COUNT(*) FILTER (WHERE status = 'canceled') AS canceled + FROM v2_job_completed + JOIN v2_job USING (id) + WHERE v2_job.workspace_id = 'admins' AND v2_job.parent_job IS NULL", + ) + .fetch_one(db) + .await + .expect("benchmark verify query failed"); + + let succeeded = row.succeeded.unwrap_or(0); + let failed = row.failed.unwrap_or(0); + let canceled = row.canceled.unwrap_or(0); + let total = succeeded + failed + canceled; + + let remaining_in_queue = sqlx::query_scalar!( + "SELECT COUNT(*) FROM v2_job_queue WHERE workspace_id = 'admins'", + ) + .fetch_one(db) + .await + .expect("benchmark verify queue query failed") + .unwrap_or(0); + + println!("=== BENCHMARK VERIFICATION ==="); + println!(" kind: {benchmark_kind}"); + println!(" expected top-level: {expected_top_level}"); + println!(" completed total: {total}"); + println!(" succeeded: {succeeded}"); + println!(" failed: {failed}"); + println!(" canceled: {canceled}"); + println!(" still in queue: {remaining_in_queue}"); + + if failed > 0 || canceled > 0 { + tracing::error!( + "BENCHMARK VERIFICATION FAILED: {failed} failed, {canceled} canceled out of {total} completed" + ); + } + if remaining_in_queue > 0 { + tracing::warn!( + "BENCHMARK VERIFICATION: {remaining_in_queue} jobs still in queue after benchmark" + ); + } + if succeeded != expected_top_level { + tracing::error!( + "BENCHMARK VERIFICATION FAILED: expected {expected_top_level} succeeded top-level jobs, got {succeeded}" + ); + } else if failed == 0 && canceled == 0 && remaining_in_queue == 0 { + println!(" result: ALL PASSED"); + } + println!("=============================="); +} + pub async fn benchmark_init(benchmark_jobs: i32, db: &DB) { use crate::{jobs::JobKind, scripts::ScriptLang}; + // Only the first worker to reach this point runs init + if BENCHMARK_INITIALIZED.swap(true, Ordering::SeqCst) { + return; + } + let benchmark_kind = std::env::var("BENCHMARK_KIND").unwrap_or("noop".to_string()); if benchmark_jobs > 0 { + // Clean up data from previous benchmark runs + sqlx::query!("DELETE FROM v2_job_completed WHERE workspace_id = 'admins'") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up v2_job_completed: {e:#}")); + sqlx::query!("DELETE FROM v2_job_queue WHERE workspace_id = 'admins'") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up v2_job_queue: {e:#}")); + sqlx::query!("DELETE FROM v2_job_status WHERE id IN (SELECT id FROM v2_job WHERE workspace_id = 'admins')") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up v2_job_status: {e:#}")); + sqlx::query!("DELETE FROM v2_job_runtime WHERE id IN (SELECT id FROM v2_job WHERE workspace_id = 'admins')") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up v2_job_runtime: {e:#}")); + sqlx::query("DELETE FROM job_perms WHERE workspace_id = 'admins'") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up job_perms: {e:#}")); + sqlx::query!("DELETE FROM concurrency_key WHERE key LIKE 'bench_%' OR key LIKE 'u/admin/bench_%'") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up concurrency_key: {e:#}")); + sqlx::query!("DELETE FROM concurrency_counter WHERE concurrency_id LIKE 'bench_%' OR concurrency_id LIKE 'u/admin/bench_%'") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up concurrency_counter: {e:#}")); + sqlx::query!("DELETE FROM v2_job WHERE workspace_id = 'admins'") + .execute(db) + .await + .unwrap_or_else(|e| panic!("failed to clean up v2_job: {e:#}")); + let mut tx = db.begin().await.unwrap(); match benchmark_kind.as_str() { "dedicated" => { @@ -260,6 +397,477 @@ pub async fn benchmark_init(benchmark_jobs: i32, db: &DB) { .await .unwrap_or_else(|_e| panic!("failed to insert parallelflow jobs (4)")); } + "sequentialflow" => { + let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_flow) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + None::, + None::, + JobKind::FlowPreview as JobKind, + ScriptLang::Deno as ScriptLang, + "flow", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + serde_json::from_str::(r#" +{ + "modules": [ + { + "id": "a", + "value": { + "type": "rawscript", + "content": "export async function main() { return 'a'; }", + "language": "deno", + "input_transforms": {} + } + }, + { + "id": "b", + "value": { + "type": "rawscript", + "content": "export async function main() { return 'b'; }", + "language": "deno", + "input_transforms": {} + } + }, + { + "id": "c", + "value": { + "type": "rawscript", + "content": "export async function main() { return 'c'; }", + "language": "deno", + "input_transforms": {} + } + }, + { + "id": "d", + "value": { + "type": "rawscript", + "content": "export async function main() { return 'd'; }", + "language": "deno", + "input_transforms": {} + } + } + ], + "preprocessor_module": null +} + "#).unwrap(), + benchmark_jobs, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (1)")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "flow") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (2)")); + sqlx::query!( + "INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", + &uuids + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (3)")); + sqlx::query!( + "INSERT INTO v2_job_status (id, flow_status) SELECT unnest($1::uuid[]), $2", + &uuids, + serde_json::from_str::( + r#" +{ + "step": 0, + "modules": [ + { "id": "a", "type": "WaitingForPriorSteps" }, + { "id": "b", "type": "WaitingForPriorSteps" }, + { "id": "c", "type": "WaitingForPriorSteps" }, + { "id": "d", "type": "WaitingForPriorSteps" } + ], + "cleanup_module": {}, + "failure_module": { + "id": "failure", + "type": "WaitingForPriorSteps" + }, + "preprocessor_module": null +} + "# + ) + .unwrap() + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (4)")); + } + "scriptlogs" => { + let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + None::, + None::, + JobKind::Preview as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + "export async function main() { for (let i = 0; i < 1000; i++) { console.log('benchmark log line ' + i); } return 'done'; }", + benchmark_jobs, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert scriptlogs jobs (1)")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert scriptlogs jobs (2)")); + sqlx::query!( + "INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", + &uuids + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert scriptlogs jobs (3)")); + } + "concurrencylimit" => { + let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id", + None::, + Some("u/admin/bench_conclimit"), + JobKind::Preview as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + "export async function main() { return 'done'; }", + 2i32, + 0i32, + benchmark_jobs, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert concurrencylimit jobs (1)")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert concurrencylimit jobs (2)")); + sqlx::query!( + "INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", + &uuids + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert concurrencylimit jobs (3)")); + let concurrency_id = "u/admin/bench_conclimit"; + sqlx::query!( + "INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING", + concurrency_id, + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert concurrencylimit counter")); + for uuid in &uuids { + sqlx::query!( + "INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)", + concurrency_id, + uuid, + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert concurrencylimit key")); + } + } + "concurrencykey" => { + let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id", + None::, + Some("u/admin/bench_conckey"), + JobKind::Preview as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + "export async function main() { return 'done'; }", + 1i32, + 0i32, + benchmark_jobs, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert concurrencykey jobs (1)")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert concurrencykey jobs (2)")); + sqlx::query!( + "INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", + &uuids + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert concurrencykey jobs (3)")); + let concurrency_id = "bench_shared_concurrency_key"; + sqlx::query!( + "INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING", + concurrency_id, + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert concurrencykey counter")); + for uuid in &uuids { + sqlx::query!( + "INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)", + concurrency_id, + uuid, + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|_e| panic!("failed to insert concurrencykey key")); + } + } + "mixed" => { + let portion = benchmark_jobs / 5; + let remainder = benchmark_jobs % 5; + + // 1) noop jobs + let noop_count = portion + remainder; + let noop_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9 FROM generate_series(1, $10)) RETURNING id", + None::, + None::, + JobKind::Noop as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + noop_count, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed noop jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &noop_uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed noop queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &noop_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed noop runtime")); + + // 2) sequentialflow jobs + if portion > 0 { + let sf_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_flow) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + None::, + None::, + JobKind::FlowPreview as JobKind, + ScriptLang::Deno as ScriptLang, + "flow", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + serde_json::from_str::(r#"{"modules":[{"id":"a","value":{"type":"rawscript","content":"export async function main() { return 'a'; }","language":"deno","input_transforms":{}}},{"id":"b","value":{"type":"rawscript","content":"export async function main() { return 'b'; }","language":"deno","input_transforms":{}}},{"id":"c","value":{"type":"rawscript","content":"export async function main() { return 'c'; }","language":"deno","input_transforms":{}}},{"id":"d","value":{"type":"rawscript","content":"export async function main() { return 'd'; }","language":"deno","input_transforms":{}}}],"preprocessor_module":null}"#).unwrap(), + portion, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sf_uuids, "admins", "flow") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sf_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow runtime")); + sqlx::query!( + "INSERT INTO v2_job_status (id, flow_status) SELECT unnest($1::uuid[]), $2", + &sf_uuids, + serde_json::from_str::(r#"{"step":0,"modules":[{"id":"a","type":"WaitingForPriorSteps"},{"id":"b","type":"WaitingForPriorSteps"},{"id":"c","type":"WaitingForPriorSteps"},{"id":"d","type":"WaitingForPriorSteps"}],"cleanup_module":{},"failure_module":{"id":"failure","type":"WaitingForPriorSteps"},"preprocessor_module":null}"#).unwrap() + ) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow status")); + } + + // 3) scriptlogs jobs + if portion > 0 { + let sl_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + None::, + None::, + JobKind::Preview as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + "export async function main() { for (let i = 0; i < 1000; i++) { console.log('benchmark log line ' + i); } return 'done'; }", + portion, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed scriptlogs jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sl_uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed scriptlogs queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sl_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed scriptlogs runtime")); + } + + // 4) concurrencylimit jobs + if portion > 0 { + let cl_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id", + None::, + Some("u/admin/bench_conclimit"), + JobKind::Preview as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + "export async function main() { return 'done'; }", + 2i32, + 0i32, + portion, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &cl_uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &cl_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit runtime")); + let cl_concurrency_id = "u/admin/bench_conclimit"; + sqlx::query!( + "INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING", + cl_concurrency_id, + ) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit counter")); + for uuid in &cl_uuids { + sqlx::query!( + "INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)", + cl_concurrency_id, + uuid, + ) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit key")); + } + } + + // 5) concurrencykey jobs + if portion > 0 { + let ck_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id", + None::, + Some("u/admin/bench_conckey"), + JobKind::Preview as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + "export async function main() { return 'done'; }", + 1i32, + 0i32, + portion, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &ck_uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &ck_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey runtime")); + let ck_concurrency_id = "bench_shared_concurrency_key"; + sqlx::query!( + "INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING", + ck_concurrency_id, + ) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey counter")); + for uuid in &ck_uuids { + sqlx::query!( + "INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)", + ck_concurrency_id, + uuid, + ) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey key")); + } + } + } + "mixed_no_cc" => { + let portion = benchmark_jobs / 3; + let remainder = benchmark_jobs % 3; + + // 1) noop jobs + let noop_count = portion + remainder; + let noop_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9 FROM generate_series(1, $10)) RETURNING id", + None::, + None::, + JobKind::Noop as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + noop_count, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc noop jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &noop_uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc noop queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &noop_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc noop runtime")); + + // 2) sequentialflow jobs + if portion > 0 { + let sf_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_flow) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + None::, + None::, + JobKind::FlowPreview as JobKind, + ScriptLang::Deno as ScriptLang, + "flow", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + serde_json::from_str::(r#"{"modules":[{"id":"a","value":{"type":"rawscript","content":"export async function main() { return 'a'; }","language":"deno","input_transforms":{}}},{"id":"b","value":{"type":"rawscript","content":"export async function main() { return 'b'; }","language":"deno","input_transforms":{}}},{"id":"c","value":{"type":"rawscript","content":"export async function main() { return 'c'; }","language":"deno","input_transforms":{}}},{"id":"d","value":{"type":"rawscript","content":"export async function main() { return 'd'; }","language":"deno","input_transforms":{}}}],"preprocessor_module":null}"#).unwrap(), + portion, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sf_uuids, "admins", "flow") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sf_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow runtime")); + sqlx::query!( + "INSERT INTO v2_job_status (id, flow_status) SELECT unnest($1::uuid[]), $2", + &sf_uuids, + serde_json::from_str::(r#"{"step":0,"modules":[{"id":"a","type":"WaitingForPriorSteps"},{"id":"b","type":"WaitingForPriorSteps"},{"id":"c","type":"WaitingForPriorSteps"},{"id":"d","type":"WaitingForPriorSteps"}],"cleanup_module":{},"failure_module":{"id":"failure","type":"WaitingForPriorSteps"},"preprocessor_module":null}"#).unwrap() + ) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow status")); + } + + // 3) scriptlogs jobs + if portion > 0 { + let sl_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + None::, + None::, + JobKind::Preview as JobKind, + ScriptLang::Deno as ScriptLang, + "deno", + "admin", + "u/admin", + "admin@windmill.dev", + "admins", + "export async function main() { for (let i = 0; i < 1000; i++) { console.log('benchmark log line ' + i); } return 'done'; }", + portion, + ) + .fetch_all(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc scriptlogs jobs")); + sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sl_uuids, "admins", "deno") + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc scriptlogs queue")); + sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sl_uuids) + .execute(&mut *tx) + .await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc scriptlogs runtime")); + } + } "none" => {} _ => { let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9 FROM generate_series(1, $10)) RETURNING id", @@ -288,6 +896,16 @@ pub async fn benchmark_init(benchmark_jobs: i32, db: &DB) { .unwrap_or_else(|_e| panic!("failed to insert noop jobs (3)")); } } + // Insert job_perms for all benchmark jobs so workers don't fall back to slow permission lookups + sqlx::query( + "INSERT INTO job_perms (job_id, email, username, is_admin, is_operator, groups, folders, workspace_id) + SELECT id, 'admin@windmill.dev', 'admin', true, false, ARRAY['all']::text[], ARRAY[]::jsonb[], 'admins' + FROM v2_job WHERE workspace_id = 'admins'" + ) + .execute(&mut *tx) + .await + .unwrap_or_else(|e| panic!("failed to insert job_perms: {e:#}")); + tx.commit().await.unwrap(); } } diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 99fa1c46ca..d707f3f8e3 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -664,153 +664,6 @@ pub async fn get_database_url() -> Result { Ok(database_url.clone()) } -pub async fn initial_connection() -> Result, error::Error> { - let connect_options = get_database_url().await?.connect_options().await?; - sqlx::postgres::PgPoolOptions::new() - .max_connections(2) - .connect_with(connect_options) - .await - .map_err(|err| Error::ConnectingToDatabase(err.to_string())) -} - -pub async fn connect_db( - server_mode: bool, - indexer_mode: bool, - worker_mode: bool, - #[cfg(feature = "private")] mut killpill_rx: tokio::sync::broadcast::Receiver<()>, -) -> anyhow::Result> { - use anyhow::Context; - - let database_url = get_database_url().await?; - - let max_connections = match std::env::var("DATABASE_CONNECTIONS") { - Ok(n) => n.parse::().context("invalid DATABASE_CONNECTIONS")?, - Err(_) => { - if server_mode { - DEFAULT_MAX_CONNECTIONS_SERVER - } else if indexer_mode { - DEFAULT_MAX_CONNECTIONS_INDEXER - } else { - DEFAULT_MAX_CONNECTIONS_WORKER - + std::env::var("NUM_WORKERS") - .ok() - .map(|x| x.parse().ok()) - .flatten() - .unwrap_or(1) - - 1 - } - } - }; - - let pool = connect(database_url.clone(), max_connections, worker_mode).await?; - #[cfg(all(feature = "enterprise", feature = "private"))] - let pool2 = pool.clone(); - #[cfg(all(feature = "enterprise", feature = "private"))] - if let DatabaseUrl::IamRds(database_url) = database_url { - tokio::spawn(async move { - loop { - tokio::select! { - _ = killpill_rx.recv() => { - break; - } - _ = tokio::time::sleep(std::time::Duration::from_secs(10)) => { - let needs_refresh = { - let read_guard = database_url.read().await; - read_guard.needs_refresh() - }; - if needs_refresh { - let new_url = tokio::time::timeout(std::time::Duration::from_secs(10), get_database_url()).await; - match new_url { - Ok(Ok(new_url)) => { - match new_url.connect_options().await { - Ok(connect_options) => { - pool2.set_connect_options(connect_options); - tracing::info!("Refreshed IAM RDS URL successfully"); - } - Err(e) => { - tracing::error!("Error getting IAM RDS connect options, retrying in 10s: {}", e); - continue; - } - } - } - Ok(Err(e)) => { - tracing::error!("Error refreshing IAM RDS URL, trying again in 10s: {}", e); - continue; - } - Err(e) => { - tracing::error!("Timeout after 10s refreshing IAM RDS URL, trying again in 10 seconds: {}", e); - continue; - } - } - } - } - } - } - }); - } - - Ok(pool) -} - -pub async fn connect( - database_url: DatabaseUrl, - max_connections: u32, - worker_mode: bool, -) -> Result, error::Error> { - use sqlx::Executor; - use std::time::Duration; - sqlx::postgres::PgPoolOptions::new() - .min_connections((max_connections / 5).clamp(3, max_connections)) - .max_connections(max_connections) - .max_lifetime(Duration::from_secs(30 * 60)) // 30 mins - .after_connect(move |conn, _| { - if worker_mode { - Box::pin(async move { - if let Err(e) = conn - .execute( - r#" - SET enable_seqscan = OFF; - SET statement_timeout = '5min'; - SET idle_in_transaction_session_timeout = '10min'; - SET tcp_keepalives_idle = 300; - SET tcp_keepalives_interval = 60; - SET tcp_keepalives_count = 10;"#, - ) - .await - { - tracing::error!("Error setting postgres settings: {}", e); - } - Ok(()) - }) - } else { - Box::pin(async move { - if let Err(e) = conn - .execute( - r#" - SET statement_timeout = '5min'; - SET idle_in_transaction_session_timeout = '10min'; - SET tcp_keepalives_idle = 300; - SET tcp_keepalives_interval = 60; - SET tcp_keepalives_count = 10;"#, - ) - .await - { - tracing::error!("Error setting postgres settings: {}", e); - } - Ok(()) - }) - } - }) - .connect_with( - database_url - .connect_options() - .await? - .statement_cache_capacity(400), - ) - .await - .map_err(|err| Error::ConnectingToDatabase(err.to_string())) -} - type Tag = String; pub use db::DB; @@ -852,11 +705,9 @@ impl ScriptHashInfo { self, db: &DB, ) -> error::Result> { - let rs = runnable_settings::from_handle( - self.runnable_settings.runnable_settings_handle, - db, - ) - .await?; + let rs = + runnable_settings::from_handle(self.runnable_settings.runnable_settings_handle, db) + .await?; let (debouncing_settings, concurrency_settings) = runnable_settings::prefetch_cached(&rs, db).await?; @@ -1176,25 +1027,25 @@ pub fn get_flow_version_info_from_version< _ => { tracing::debug!("Fetching flow version info for {version} ({path})"); let mut conn = db.acquire().await?; - let flow_info = + let flow_info = sqlx::query_as!( FlowVersionInfo, r#" SELECT flow_version.id AS version, - flow_version.value->>'early_return' as early_return, - flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor, - (flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled, - flow.tag, - flow.dedicated_worker, - flow.on_behalf_of_email, + flow_version.value->>'early_return' as early_return, + flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor, + (flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled, + flow.tag, + flow.dedicated_worker, + flow.on_behalf_of_email, flow.edited_by - FROM + FROM flow_version INNER JOIN flow ON flow.path = flow_version.path AND flow.workspace_id = flow_version.workspace_id - WHERE + WHERE flow_version.workspace_id = $1 AND flow_version.path = $2 AND flow_version.id = $3 diff --git a/backend/windmill-common/src/utils.rs b/backend/windmill-common/src/utils.rs index 4a6f4604a1..9f6bbf7893 100644 --- a/backend/windmill-common/src/utils.rs +++ b/backend/windmill-common/src/utils.rs @@ -145,6 +145,13 @@ lazy_static::lazy_static! { tracing::info!("Mode not specified, defaulting to standalone"); Mode::Standalone }); + #[cfg(feature = "benchmark")] + let mode = { + if mode != Mode::Worker { + println!("Benchmark mode: forcing MODE=worker"); + } + Mode::Worker + }; ModeAndAddons { indexer: search_addon, mode, diff --git a/backend/windmill-worker/src/result_processor.rs b/backend/windmill-worker/src/result_processor.rs index 6992d6ff74..c237d8cc02 100644 --- a/backend/windmill-worker/src/result_processor.rs +++ b/backend/windmill-worker/src/result_processor.rs @@ -68,6 +68,7 @@ async fn process_jc( stats_map: &JobStatsMap, killpill_rx: &tokio::sync::broadcast::Receiver<()>, #[cfg(feature = "benchmark")] bench: &mut BenchmarkIter, + #[cfg(feature = "benchmark")] bench_infos: &mut BenchmarkInfo, ) { let success: bool = jc.success; @@ -166,6 +167,13 @@ async fn process_jc( if let Some(root_job) = root_job { add_root_flow_job_to_otlp(&root_job, success); + + #[cfg(feature = "benchmark")] + if bench_infos.count_top_level(root_job.id) { + bench_infos + .shared_iters + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + } } // Accumulate job stats if duration is available @@ -209,7 +217,7 @@ pub fn start_background_processor( let JobCompletedReceiver { bounded_rx, mut killpill_rx, unbounded_rx } = job_completed_rx; #[cfg(feature = "benchmark")] - let mut infos = BenchmarkInfo::new(); + let mut infos = BenchmarkInfo::new(windmill_common::bench::shared_bench_iters()); // Start periodic stats flush task let db_clone = db.clone(); @@ -278,6 +286,10 @@ pub fn start_background_processor( jc.job.kind, JobKind::Dependencies | JobKind::FlowDependencies ); + #[cfg(feature = "benchmark")] + let bench_job_id = jc.job.id; + #[cfg(feature = "benchmark")] + let is_top_level_job = jc.job.parent_job.is_none(); process_jc( jc, @@ -291,6 +303,8 @@ pub fn start_background_processor( &killpill_rx, #[cfg(feature = "benchmark")] &mut bench, + #[cfg(feature = "benchmark")] + &mut infos, ) .warn_after_seconds(10) .await; @@ -315,7 +329,9 @@ pub fn start_background_processor( #[cfg(feature = "benchmark")] { - infos.add_iter(bench, true); + if infos.add_iter(bench, bench_job_id, is_top_level_job) { + infos.shared_iters.fetch_add(1, Ordering::Relaxed); + } } last_processing_duration .store(time.elapsed().as_secs() as u16, Ordering::SeqCst); @@ -365,6 +381,12 @@ pub fn start_background_processor( { tracing::error!("Error updating flow status after job completion for {flow} on {worker_name}: {e:#}"); } + #[cfg(feature = "benchmark")] + { + if infos.add_iter(bench, flow, true) { + infos.shared_iters.fetch_add(1, Ordering::Relaxed); + } + } last_processing_duration .store(time.elapsed().as_secs() as u16, Ordering::SeqCst); } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index f3cdc241e0..06f59e3814 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -188,7 +188,7 @@ use crate::mssql_executor::do_mssql; use crate::bigquery_executor::do_bigquery; #[cfg(feature = "benchmark")] -use windmill_common::bench::{benchmark_init, BenchmarkInfo, BenchmarkIter}; +use windmill_common::bench::{benchmark_init, benchmark_verify, BenchmarkInfo, BenchmarkIter}; use windmill_common::add_time; @@ -1109,51 +1109,60 @@ fn start_interactive_worker_shell( if let Ok(_) = killpill_rx.try_recv() { tracing::info!("Received killpill, exiting worker shell"); break; - } else { - let pulled_job = match &conn { - Connection::Sql(db) => { - let common_worker_prefix = retrieve_common_worker_prefix(&worker_name); - let query = ("".to_string(), make_pull_query(&[common_worker_prefix])); - #[cfg(feature = "benchmark")] - let mut bench = windmill_common::bench::BenchmarkIter::new(); + } - let job = pull( - &db, - false, - &worker_name, - Some(&query), + let pulled_job = tokio::select! { + _ = killpill_rx.recv() => { + tracing::info!("Received killpill during pull, exiting worker shell"); + break; + } + result = async { + match &conn { + Connection::Sql(db) => { + let common_worker_prefix = retrieve_common_worker_prefix(&worker_name); + let query = ("".to_string(), make_pull_query(&[common_worker_prefix])); #[cfg(feature = "benchmark")] - &mut bench, - ) - .await; + let mut bench = windmill_common::bench::BenchmarkIter::new(); - use PulledJobResultToJobErr::*; - match job { - Ok(j) => match j.to_pulled_job() { - Ok(j) => Ok(j - .clone() - .map(|job| NextJob::Sql { flow_runners: None, job })), - Err(MissingConcurrencyKey(jc)) - | Err(ErrorWhilePreprocessing(jc)) => { - if let Err(err) = job_completed_tx.send_job(jc, true).await { - tracing::error!( - "An error occurred while sending job completed: {:#?}", - err - ) + let job = pull( + &db, + false, + &worker_name, + Some(&query), + #[cfg(feature = "benchmark")] + &mut bench, + ) + .await; + + use PulledJobResultToJobErr::*; + match job { + Ok(j) => match j.to_pulled_job() { + Ok(j) => Ok(j + .clone() + .map(|job| NextJob::Sql { flow_runners: None, job })), + Err(MissingConcurrencyKey(jc)) + | Err(ErrorWhilePreprocessing(jc)) => { + if let Err(err) = job_completed_tx.send_job(jc, true).await { + tracing::error!( + "An error occurred while sending job completed: {:#?}", + err + ) + } + Ok(None) } - Ok(None) - } - }, - Err(err) => Err(err), + }, + Err(err) => Err(err), + } + } + Connection::Http(client) => { + crate::agent_workers::pull_job(&client, None, Some(true)) + .await + .map_err(|e| error::Error::InternalErr(e.to_string())) + .map(|x| x.map(|y| NextJob::Http(y))) } } - Connection::Http(client) => { - crate::agent_workers::pull_job(&client, None, Some(true)) - .await - .map_err(|e| error::Error::InternalErr(e.to_string())) - .map(|x| x.map(|y| NextJob::Http(y))) - } - }; + } => result, + }; match pulled_job { Ok(Some(job)) => { @@ -1233,7 +1242,6 @@ fn start_interactive_worker_shell( tokio::time::sleep(Duration::from_millis(*SLEEP_QUEUE * 20)).await; } }; - } } }) } @@ -1571,6 +1579,7 @@ pub async fn run_worker( // This is used to wake up the background processor when main loop is done and just waiting for new same workers jobs, and that bg processor is also not processing any jobs, bg processing can exit if no more same worker jobs let wake_up_notify = Arc::new(tokio::sync::Notify::new()); let stats_map = JobStatsMap::default(); + let send_result = match (conn, job_completed_rx) { (Connection::Sql(db), Some(job_completed_receiver)) => Some(start_background_processor( job_completed_receiver, @@ -1615,7 +1624,10 @@ pub async fn run_worker( let mut started = false; #[cfg(feature = "benchmark")] - let mut infos = BenchmarkInfo::new(); + let mut infos = BenchmarkInfo::new(windmill_common::bench::shared_bench_iters()); + + #[cfg(feature = "benchmark")] + let mut bench_empty_queue_count: u64 = 0; #[cfg(feature = "benchmark")] if let Some(db) = conn.as_sql() { @@ -1807,15 +1819,60 @@ pub async fn run_worker( // } #[cfg(feature = "benchmark")] - if benchmark_jobs > 0 && infos.iters == benchmark_jobs as u64 { - tracing::info!("benchmark finished, exiting"); - job_completed_tx - .kill() - .await - .expect("send kill to job completed tx"); - break; - } else { - tracing::info!("benchmark not finished, still pulling jobs {}", infos.iters); + { + let total_iters = infos.shared_iters.load(std::sync::atomic::Ordering::Relaxed); + if benchmark_jobs > 0 && total_iters >= benchmark_jobs as u64 { + tracing::info!("benchmark finished, exiting (total iters: {}, worker iters: {})", total_iters, infos.iters); + job_completed_tx + .kill() + .await + .expect("send kill to job completed tx"); + killpill_tx.send(); + break; + } else if benchmark_jobs > 0 && bench_empty_queue_count > 2000 { + tracing::warn!( + "benchmark stalled: no jobs in queue for 2000 polls, exiting (total iters: {}, worker iters: {}/{})", + total_iters, + infos.iters, + benchmark_jobs + ); + job_completed_tx + .kill() + .await + .expect("send kill to job completed tx"); + killpill_tx.send(); + break; + } else if bench_empty_queue_count % 100 == 0 { + if let Some(db) = conn.as_sql() { + let remaining = sqlx::query_as::<_, (uuid::Uuid, String, bool, Option, Option)>( + "SELECT q.id, q.tag, q.running, j.kind::text, j.parent_job + FROM v2_job_queue q JOIN v2_job j ON q.id = j.id + WHERE q.workspace_id = 'admins' LIMIT 10" + ) + .fetch_all(db) + .await; + match remaining { + Ok(rows) => { + let total_remaining = sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM v2_job_queue WHERE workspace_id = 'admins'" + ).fetch_one(db).await.unwrap_or(0); + for (id, tag, running, kind, parent) in &rows { + tracing::info!( + " pending job: id={id}, tag={tag}, running={running}, kind={}, parent={:?}", + kind.as_deref().unwrap_or("?"), parent + ); + } + tracing::info!( + "benchmark not finished (total: {}, worker: {}, queue: {})", + total_iters, infos.iters, total_remaining + ); + } + Err(e) => { + tracing::info!("benchmark not finished (total: {}, worker: {}), queue query err: {e}", total_iters, infos.iters); + } + } + } + } } let next_job = { @@ -2029,6 +2086,15 @@ pub async fn run_worker( match next_job { Ok(Some(job)) => { + #[cfg(feature = "benchmark")] + { + bench_empty_queue_count = 0; + } + #[cfg(feature = "benchmark")] + let is_top_level_job = job.parent_job.is_none() && !job.kind.is_flow(); + #[cfg(feature = "benchmark")] + let bench_job_id = job.id; + #[cfg(feature = "prometheus")] if let Some(wb) = worker_busy.as_ref() { wb.set(1); @@ -2074,7 +2140,7 @@ pub async fn run_worker( if let Some(db) = conn.as_sql() { infos.sample_pool(db.size(), db.num_idle() as u32); } - infos.add_iter(bench, true); + infos.add_iter(bench, bench_job_id, is_top_level_job); } continue; @@ -2147,7 +2213,7 @@ pub async fn run_worker( if let Some(db) = conn.as_sql() { infos.sample_pool(db.size(), db.num_idle() as u32); } - infos.add_iter(bench, true); + infos.add_iter(bench, bench_job_id, is_top_level_job); } continue; @@ -2450,7 +2516,7 @@ pub async fn run_worker( if let Some(db) = conn.as_sql() { infos.sample_pool(db.size(), db.num_idle() as u32); } - infos.add_iter(bench, true); + infos.add_iter(bench, bench_job_id, is_top_level_job); } } } @@ -2477,11 +2543,12 @@ pub async fn run_worker( #[cfg(feature = "benchmark")] { + bench_empty_queue_count += 1; add_time!(bench, "sleep because empty job queue"); if let Some(db) = conn.as_sql() { infos.sample_pool(db.size(), db.num_idle() as u32); } - infos.add_iter(bench, false); + infos.add_iter(bench, uuid::Uuid::nil(), false); } #[cfg(feature = "prometheus")] _timer.map(|timer| { @@ -2510,13 +2577,6 @@ pub async fn run_worker( } } - #[cfg(feature = "benchmark")] - { - infos - .write_to_file("profiling_main.json") - .expect("write to file profiling"); - } - drop(dedicated_workers); let has_dedicated_workers = !dedicated_handles.is_empty(); @@ -2537,6 +2597,17 @@ pub async fn run_worker( tracing::error!("error in awaiting send_result process: {e:?}") } } + + #[cfg(feature = "benchmark")] + { + infos + .write_to_file("profiling_main.json") + .expect("write to file profiling"); + + if let Some(db) = conn.as_sql() { + benchmark_verify(benchmark_jobs, db).await; + } + } tracing::info!(worker = %worker_name, hostname = %hostname, "waiting for interactive_shell to finish"); if let Some(interactive_shell) = interactive_shell { match tokio::time::timeout(Duration::from_secs(10), interactive_shell).await { diff --git a/frontend/src/lib/components/DeployWorkspace.svelte b/frontend/src/lib/components/DeployWorkspace.svelte index 0624240bd0..c7c5ba2ca0 100644 --- a/frontend/src/lib/components/DeployWorkspace.svelte +++ b/frontend/src/lib/components/DeployWorkspace.svelte @@ -149,7 +149,6 @@ }) } else if (kind == 'app') { const app = await AppService.getAppByPath({ workspace: $workspaceStore!, path }) - console.log('app', app) let result: { kind: Kind; path: string }[] = [] if (app.raw_app) { const rawAppValue = app.value as { runnables?: Record }