diff --git a/backend/.sqlx/query-07d03985bb2c58d52c1ffd6ab5a6d37457e7520642a5e70bb4000e4923720957.json b/backend/.sqlx/query-07d03985bb2c58d52c1ffd6ab5a6d37457e7520642a5e70bb4000e4923720957.json new file mode 100644 index 0000000000..c8bc7af4d3 --- /dev/null +++ b/backend/.sqlx/query-07d03985bb2c58d52c1ffd6ab5a6d37457e7520642a5e70bb4000e4923720957.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE flow_version SET path = REGEXP_REPLACE(path,'u/' || $2 || '/(.*)','u/' || $1 || '/\\1') WHERE path LIKE ('u/' || $2 || '/%') AND workspace_id = $3", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "07d03985bb2c58d52c1ffd6ab5a6d37457e7520642a5e70bb4000e4923720957" +} diff --git a/backend/.sqlx/query-07f5290e90533eac50b890a0d7f4a5e73ac111c838f687fe8647636827aae8b5.json b/backend/.sqlx/query-07f5290e90533eac50b890a0d7f4a5e73ac111c838f687fe8647636827aae8b5.json new file mode 100644 index 0000000000..811920b354 --- /dev/null +++ b/backend/.sqlx/query-07f5290e90533eac50b890a0d7f4a5e73ac111c838f687fe8647636827aae8b5.json @@ -0,0 +1,26 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO flow_version (workspace_id, path, value, schema, created_by) \n VALUES ($1, $2, $3, $4::text::json, $5)\n RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Jsonb", + "Text", + "Varchar" + ] + }, + "nullable": [ + false + ] + }, + "hash": "07f5290e90533eac50b890a0d7f4a5e73ac111c838f687fe8647636827aae8b5" +} diff --git a/backend/.sqlx/query-13358ffeb0917dd9dff9f8527a59dfee63bb704c3f712af179732dc281411917.json b/backend/.sqlx/query-13358ffeb0917dd9dff9f8527a59dfee63bb704c3f712af179732dc281411917.json new file mode 100644 index 0000000000..d49a1b9d44 --- /dev/null +++ b/backend/.sqlx/query-13358ffeb0917dd9dff9f8527a59dfee63bb704c3f712af179732dc281411917.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO flow \n (workspace_id, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at) \n SELECT $1, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at\n FROM flow WHERE workspace_id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Text" + ] + }, + "nullable": [] + }, + "hash": "13358ffeb0917dd9dff9f8527a59dfee63bb704c3f712af179732dc281411917" +} diff --git a/backend/.sqlx/query-25155e44372aecbb38d042bfc2772ed0c01a0bb974488530cd713b834f537f4a.json b/backend/.sqlx/query-25155e44372aecbb38d042bfc2772ed0c01a0bb974488530cd713b834f537f4a.json new file mode 100644 index 0000000000..f459a72957 --- /dev/null +++ b/backend/.sqlx/query-25155e44372aecbb38d042bfc2772ed0c01a0bb974488530cd713b834f537f4a.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO flow\n (workspace_id, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at) \n SELECT workspace_id, REGEXP_REPLACE(path,'u/' || $2 || '/(.*)','u/' || $1 || '/\\1'), summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at\n FROM flow \n WHERE path LIKE ('u/' || $2 || '/%') AND workspace_id = $3", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "25155e44372aecbb38d042bfc2772ed0c01a0bb974488530cd713b834f537f4a" +} diff --git a/backend/.sqlx/query-34dee810f99ef41727ab3231a1746be80d60050f8cbaf779d391c4e08eb0c438.json b/backend/.sqlx/query-34dee810f99ef41727ab3231a1746be80d60050f8cbaf779d391c4e08eb0c438.json new file mode 100644 index 0000000000..2b8c9373e6 --- /dev/null +++ b/backend/.sqlx/query-34dee810f99ef41727ab3231a1746be80d60050f8cbaf779d391c4e08eb0c438.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE flow_version SET value = $1 WHERE id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Jsonb", + "Int8" + ] + }, + "nullable": [] + }, + "hash": "34dee810f99ef41727ab3231a1746be80d60050f8cbaf779d391c4e08eb0c438" +} diff --git a/backend/.sqlx/query-729bf764b10dc821594bbffbc157c061ec6eb09b3cc22c672e0e17da455ac7ef.json b/backend/.sqlx/query-4502ed44e69b0501ef187be9cc2de22c4dc5dafb13ef56c297ea62464e74c323.json similarity index 53% rename from backend/.sqlx/query-729bf764b10dc821594bbffbc157c061ec6eb09b3cc22c672e0e17da455ac7ef.json rename to backend/.sqlx/query-4502ed44e69b0501ef187be9cc2de22c4dc5dafb13ef56c297ea62464e74c323.json index c32ba28341..c17b642f5b 100644 --- a/backend/.sqlx/query-729bf764b10dc821594bbffbc157c061ec6eb09b3cc22c672e0e17da455ac7ef.json +++ b/backend/.sqlx/query-4502ed44e69b0501ef187be9cc2de22c4dc5dafb13ef56c297ea62464e74c323.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE flow SET edited_by = $1 WHERE edited_by = $2 AND workspace_id = $3", + "query": "UPDATE flow_version SET path = $1 WHERE path = $2 AND workspace_id = $3", "describe": { "columns": [], "parameters": { @@ -12,5 +12,5 @@ }, "nullable": [] }, - "hash": "729bf764b10dc821594bbffbc157c061ec6eb09b3cc22c672e0e17da455ac7ef" + "hash": "4502ed44e69b0501ef187be9cc2de22c4dc5dafb13ef56c297ea62464e74c323" } diff --git a/backend/.sqlx/query-7f9f1ce221835fc3ff6c864c362519fde3d08e496119e51d8ed4fe65792a85c1.json b/backend/.sqlx/query-4968e9edac534657c808b891cbf93c8c0a57f93b7b445171b1cc3f4428ee6e53.json similarity index 50% rename from backend/.sqlx/query-7f9f1ce221835fc3ff6c864c362519fde3d08e496119e51d8ed4fe65792a85c1.json rename to backend/.sqlx/query-4968e9edac534657c808b891cbf93c8c0a57f93b7b445171b1cc3f4428ee6e53.json index f328b616e0..08bc20e88b 100644 --- a/backend/.sqlx/query-7f9f1ce221835fc3ff6c864c362519fde3d08e496119e51d8ed4fe65792a85c1.json +++ b/backend/.sqlx/query-4968e9edac534657c808b891cbf93c8c0a57f93b7b445171b1cc3f4428ee6e53.json @@ -1,17 +1,18 @@ { "db_name": "PostgreSQL", - "query": "INSERT INTO deployment_metadata (workspace_id, path, callback_job_ids, deployment_msg) VALUES ($1, $2, $3, $4) ON CONFLICT (workspace_id, path) WHERE script_hash IS NULL AND app_version IS NULL DO UPDATE SET callback_job_ids = $3, deployment_msg = $4", + "query": "INSERT INTO deployment_metadata (workspace_id, path, flow_version, callback_job_ids, deployment_msg) VALUES ($1, $2, $3, $4, $5)", "describe": { "columns": [], "parameters": { "Left": [ "Varchar", "Varchar", + "Int8", "UuidArray", "Text" ] }, "nullable": [] }, - "hash": "7f9f1ce221835fc3ff6c864c362519fde3d08e496119e51d8ed4fe65792a85c1" + "hash": "4968e9edac534657c808b891cbf93c8c0a57f93b7b445171b1cc3f4428ee6e53" } diff --git a/backend/.sqlx/query-4a764cdcb847b71183425a7a0e863708ef7fe2c88b28d0667988a47b9e995c0e.json b/backend/.sqlx/query-4a764cdcb847b71183425a7a0e863708ef7fe2c88b28d0667988a47b9e995c0e.json new file mode 100644 index 0000000000..fd3f7b2b4f --- /dev/null +++ b/backend/.sqlx/query-4a764cdcb847b71183425a7a0e863708ef7fe2c88b28d0667988a47b9e995c0e.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT flow_version.value \n FROM flow \n LEFT JOIN flow_version \n ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]\n WHERE flow.path = $1 AND flow.workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "value", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + true + ] + }, + "hash": "4a764cdcb847b71183425a7a0e863708ef7fe2c88b28d0667988a47b9e995c0e" +} diff --git a/backend/.sqlx/query-526bfaccaafbe2e6f70dd5e6cd21c0c60d4ec155f79d067a8b74cf24eebad88c.json b/backend/.sqlx/query-526bfaccaafbe2e6f70dd5e6cd21c0c60d4ec155f79d067a8b74cf24eebad88c.json new file mode 100644 index 0000000000..08e38953ab --- /dev/null +++ b/backend/.sqlx/query-526bfaccaafbe2e6f70dd5e6cd21c0c60d4ec155f79d067a8b74cf24eebad88c.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT versions[array_upper(versions, 1)] FROM flow WHERE path = $1 AND workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "versions", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "526bfaccaafbe2e6f70dd5e6cd21c0c60d4ec155f79d067a8b74cf24eebad88c" +} diff --git a/backend/.sqlx/query-5ff7df54c7908a7de494ddae5fc7bb9be8106a79e0683cd34459585bbd920ce4.json b/backend/.sqlx/query-5ff7df54c7908a7de494ddae5fc7bb9be8106a79e0683cd34459585bbd920ce4.json new file mode 100644 index 0000000000..45bca31fd0 --- /dev/null +++ b/backend/.sqlx/query-5ff7df54c7908a7de494ddae5fc7bb9be8106a79e0683cd34459585bbd920ce4.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT flow_version.value->>'concurrency_key'\n FROM flow \n LEFT JOIN flow_version\n ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]\n WHERE flow.path = $1 AND flow.workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "?column?", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "5ff7df54c7908a7de494ddae5fc7bb9be8106a79e0683cd34459585bbd920ce4" +} diff --git a/backend/.sqlx/query-6b1ea6f39c6f41a093112418ac4c1b69a57de50fdeb4bdcd8d4cab0553af42ea.json b/backend/.sqlx/query-6b1ea6f39c6f41a093112418ac4c1b69a57de50fdeb4bdcd8d4cab0553af42ea.json deleted file mode 100644 index 1bc6db6cfa..0000000000 --- a/backend/.sqlx/query-6b1ea6f39c6f41a093112418ac4c1b69a57de50fdeb4bdcd8d4cab0553af42ea.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT value FROM flow WHERE path = $1 AND workspace_id = $2", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "value", - "type_info": "Jsonb" - } - ], - "parameters": { - "Left": [ - "Text", - "Text" - ] - }, - "nullable": [ - false - ] - }, - "hash": "6b1ea6f39c6f41a093112418ac4c1b69a57de50fdeb4bdcd8d4cab0553af42ea" -} diff --git a/backend/.sqlx/query-76e1de02790d23394997eeec6a5ee46d1da97b94bdd4a07af3dc57b7f8e6f089.json b/backend/.sqlx/query-76e1de02790d23394997eeec6a5ee46d1da97b94bdd4a07af3dc57b7f8e6f089.json new file mode 100644 index 0000000000..181da1762d --- /dev/null +++ b/backend/.sqlx/query-76e1de02790d23394997eeec6a5ee46d1da97b94bdd4a07af3dc57b7f8e6f089.json @@ -0,0 +1,26 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO flow_version (workspace_id, path, value, schema, created_by) VALUES ($1, $2, $3, $4::text::json, $5) RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Jsonb", + "Text", + "Varchar" + ] + }, + "nullable": [ + false + ] + }, + "hash": "76e1de02790d23394997eeec6a5ee46d1da97b94bdd4a07af3dc57b7f8e6f089" +} diff --git a/backend/.sqlx/query-79464d5ef46a05ff9c05a4f1f4419ffac7e82c985d59a2e20e5c616461dbfe7b.json b/backend/.sqlx/query-79464d5ef46a05ff9c05a4f1f4419ffac7e82c985d59a2e20e5c616461dbfe7b.json new file mode 100644 index 0000000000..1cdb56d8c0 --- /dev/null +++ b/backend/.sqlx/query-79464d5ef46a05ff9c05a4f1f4419ffac7e82c985d59a2e20e5c616461dbfe7b.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT flow.path FROM flow\n LEFT JOIN flow_version\n ON flow_version.path = flow.path AND flow_version.workspace_id = flow.workspace_id\n WHERE flow.path = $1 AND flow.workspace_id = $2 AND flow_version.id = $3", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "path", + "type_info": "Varchar" + } + ], + "parameters": { + "Left": [ + "Text", + "Text", + "Int8" + ] + }, + "nullable": [ + false + ] + }, + "hash": "79464d5ef46a05ff9c05a4f1f4419ffac7e82c985d59a2e20e5c616461dbfe7b" +} diff --git a/backend/.sqlx/query-8e0679c2b1bd451691fe5c69a2841ddc9f211311316ec6b9d4699b2c70997a19.json b/backend/.sqlx/query-872dcaec230579e4480adf23075e323557efbe52c813e3a6a0da6b855291951e.json similarity index 58% rename from backend/.sqlx/query-8e0679c2b1bd451691fe5c69a2841ddc9f211311316ec6b9d4699b2c70997a19.json rename to backend/.sqlx/query-872dcaec230579e4480adf23075e323557efbe52c813e3a6a0da6b855291951e.json index 9f2a6fabd8..cdf6a3f4fa 100644 --- a/backend/.sqlx/query-8e0679c2b1bd451691fe5c69a2841ddc9f211311316ec6b9d4699b2c70997a19.json +++ b/backend/.sqlx/query-872dcaec230579e4480adf23075e323557efbe52c813e3a6a0da6b855291951e.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT tag, dedicated_worker, value->>'early_return' as early_return from flow WHERE path = $1 and workspace_id = $2", + "query": "SELECT tag, dedicated_worker, flow_version.value->>'early_return' as early_return \n FROM flow \n LEFT JOIN flow_version\n ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]\n WHERE flow.path = $1 and flow.workspace_id = $2", "describe": { "columns": [ { @@ -31,5 +31,5 @@ null ] }, - "hash": "8e0679c2b1bd451691fe5c69a2841ddc9f211311316ec6b9d4699b2c70997a19" + "hash": "872dcaec230579e4480adf23075e323557efbe52c813e3a6a0da6b855291951e" } diff --git a/backend/.sqlx/query-f6fd65fbe36502923ab4ccf1a22f748cb854e23049d0cf73a42229eccc88c4e6.json b/backend/.sqlx/query-899a406cc03659e29bad831cc4d3d1dcda1a39b826a03136353344ccae29871b.json similarity index 51% rename from backend/.sqlx/query-f6fd65fbe36502923ab4ccf1a22f748cb854e23049d0cf73a42229eccc88c4e6.json rename to backend/.sqlx/query-899a406cc03659e29bad831cc4d3d1dcda1a39b826a03136353344ccae29871b.json index 72dff20b56..dec3a1bb5a 100644 --- a/backend/.sqlx/query-f6fd65fbe36502923ab4ccf1a22f748cb854e23049d0cf73a42229eccc88c4e6.json +++ b/backend/.sqlx/query-899a406cc03659e29bad831cc4d3d1dcda1a39b826a03136353344ccae29871b.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE flow SET path = $1, summary = $2, description = $3, value = $4, edited_by = $5, edited_at = now(), schema = $6::text::json, dependency_job = NULL, draft_only = NULL, tag = $9, dedicated_worker = $10, visible_to_runner_only = $11\n WHERE path = $7 AND workspace_id = $8", + "query": "UPDATE flow SET path = $1, summary = $2, description = $3,dependency_job = NULL, draft_only = NULL, tag = $4, dedicated_worker = $5, visible_to_runner_only = $6, value = $7, schema = $8::text::json, edited_by = $9, edited_at = now()\n WHERE path = $10 AND workspace_id = $11", "describe": { "columns": [], "parameters": { @@ -8,17 +8,17 @@ "Varchar", "Text", "Text", - "Jsonb", - "Varchar", - "Text", - "Text", - "Text", "Varchar", "Bool", - "Bool" + "Bool", + "Jsonb", + "Text", + "Varchar", + "Text", + "Text" ] }, "nullable": [] }, - "hash": "f6fd65fbe36502923ab4ccf1a22f748cb854e23049d0cf73a42229eccc88c4e6" + "hash": "899a406cc03659e29bad831cc4d3d1dcda1a39b826a03136353344ccae29871b" } diff --git a/backend/.sqlx/query-31075ff185a9ab857459bc539eadd1022c1e5bf0cfbd02c97739f5b83350f050.json b/backend/.sqlx/query-974c7e623f3dfa440e134eaaa8d029334c0645147200219c39b2c00b30941172.json similarity index 66% rename from backend/.sqlx/query-31075ff185a9ab857459bc539eadd1022c1e5bf0cfbd02c97739f5b83350f050.json rename to backend/.sqlx/query-974c7e623f3dfa440e134eaaa8d029334c0645147200219c39b2c00b30941172.json index fcb37677ab..0a735b5a2c 100644 --- a/backend/.sqlx/query-31075ff185a9ab857459bc539eadd1022c1e5bf0cfbd02c97739f5b83350f050.json +++ b/backend/.sqlx/query-974c7e623f3dfa440e134eaaa8d029334c0645147200219c39b2c00b30941172.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT workspace_id as workspace, path, summary, description, schema FROM flow WHERE workspace_id = $1", + "query": "SELECT flow.workspace_id as workspace, flow.path, summary, description, flow_version.schema \n FROM flow \n LEFT JOIN flow_version ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]\n WHERE flow.workspace_id = $1", "describe": { "columns": [ { @@ -42,5 +42,5 @@ true ] }, - "hash": "31075ff185a9ab857459bc539eadd1022c1e5bf0cfbd02c97739f5b83350f050" + "hash": "974c7e623f3dfa440e134eaaa8d029334c0645147200219c39b2c00b30941172" } diff --git a/backend/.sqlx/query-35e6af0b203e3e4fac9020b037a3c41af92537c0fd5683227767f6a1bd17339f.json b/backend/.sqlx/query-9e64c6b6db2155ad8e6e514b08da92d80caafd1a217813bcd354aa264810f690.json similarity index 52% rename from backend/.sqlx/query-35e6af0b203e3e4fac9020b037a3c41af92537c0fd5683227767f6a1bd17339f.json rename to backend/.sqlx/query-9e64c6b6db2155ad8e6e514b08da92d80caafd1a217813bcd354aa264810f690.json index 8076362b64..25296490e0 100644 --- a/backend/.sqlx/query-35e6af0b203e3e4fac9020b037a3c41af92537c0fd5683227767f6a1bd17339f.json +++ b/backend/.sqlx/query-9e64c6b6db2155ad8e6e514b08da92d80caafd1a217813bcd354aa264810f690.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "INSERT INTO flow (workspace_id, path, summary, description, value, edited_by, edited_at, schema, dependency_job, draft_only, tag, dedicated_worker, visible_to_runner_only) VALUES ($1, $2, $3, $4, $5, $6, now(), $7::text::json, NULL, $8, $9, $10, $11)", + "query": "INSERT INTO flow (workspace_id, path, summary, description, dependency_job, draft_only, tag, dedicated_worker, visible_to_runner_only, value, schema, edited_by, edited_at) \n VALUES ($1, $2, $3, $4, NULL, $5, $6, $7, $8, $9, $10::text::json, $11, now())", "describe": { "columns": [], "parameters": { @@ -9,16 +9,16 @@ "Varchar", "Text", "Text", + "Bool", + "Varchar", + "Bool", + "Bool", "Jsonb", - "Varchar", "Text", - "Bool", - "Varchar", - "Bool", - "Bool" + "Varchar" ] }, "nullable": [] }, - "hash": "35e6af0b203e3e4fac9020b037a3c41af92537c0fd5683227767f6a1bd17339f" + "hash": "9e64c6b6db2155ad8e6e514b08da92d80caafd1a217813bcd354aa264810f690" } diff --git a/backend/.sqlx/query-a0f1c0df6bc2f1fbca50edee90e42c94445536e201b322eda6f7a90bdf38f36a.json b/backend/.sqlx/query-a0f1c0df6bc2f1fbca50edee90e42c94445536e201b322eda6f7a90bdf38f36a.json new file mode 100644 index 0000000000..3ed4f5a316 --- /dev/null +++ b/backend/.sqlx/query-a0f1c0df6bc2f1fbca50edee90e42c94445536e201b322eda6f7a90bdf38f36a.json @@ -0,0 +1,35 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT flow_version.id, flow_version.created_at, deployment_metadata.deployment_msg FROM flow_version \n LEFT JOIN deployment_metadata ON flow_version.id = deployment_metadata.flow_version\n WHERE flow_version.path = $1 AND flow_version.workspace_id = $2 \n ORDER BY flow_version.created_at DESC", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "created_at", + "type_info": "Timestamptz" + }, + { + "ordinal": 2, + "name": "deployment_msg", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + false, + false, + true + ] + }, + "hash": "a0f1c0df6bc2f1fbca50edee90e42c94445536e201b322eda6f7a90bdf38f36a" +} diff --git a/backend/.sqlx/query-a459d973c6392af6af1af8b17bc735128fa06ee80e867e2421e70df7a722081a.json b/backend/.sqlx/query-a459d973c6392af6af1af8b17bc735128fa06ee80e867e2421e70df7a722081a.json deleted file mode 100644 index 7e08e3ca3c..0000000000 --- a/backend/.sqlx/query-a459d973c6392af6af1af8b17bc735128fa06ee80e867e2421e70df7a722081a.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE flow SET path = REGEXP_REPLACE(path,'u/' || $2 || '/(.*)','u/' || $1 || '/\\1') WHERE path LIKE ('u/' || $2 || '/%') AND workspace_id = $3", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Text", - "Text", - "Text" - ] - }, - "nullable": [] - }, - "hash": "a459d973c6392af6af1af8b17bc735128fa06ee80e867e2421e70df7a722081a" -} diff --git a/backend/.sqlx/query-a875cb56485b812e9d4739afd0915067f7e5abe0ca0adf264b792fccf21e005b.json b/backend/.sqlx/query-a875cb56485b812e9d4739afd0915067f7e5abe0ca0adf264b792fccf21e005b.json deleted file mode 100644 index e264db5ef2..0000000000 --- a/backend/.sqlx/query-a875cb56485b812e9d4739afd0915067f7e5abe0ca0adf264b792fccf21e005b.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT value->>'concurrency_key' FROM flow WHERE path = $1 AND workspace_id = $2", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "?column?", - "type_info": "Text" - } - ], - "parameters": { - "Left": [ - "Text", - "Text" - ] - }, - "nullable": [ - null - ] - }, - "hash": "a875cb56485b812e9d4739afd0915067f7e5abe0ca0adf264b792fccf21e005b" -} diff --git a/backend/.sqlx/query-afc7c23c057748f6d4a61dbef17e433b8875c6588b91e38e8141d12189118412.json b/backend/.sqlx/query-afc7c23c057748f6d4a61dbef17e433b8875c6588b91e38e8141d12189118412.json new file mode 100644 index 0000000000..08ea345138 --- /dev/null +++ b/backend/.sqlx/query-afc7c23c057748f6d4a61dbef17e433b8875c6588b91e38e8141d12189118412.json @@ -0,0 +1,17 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO deployment_metadata (workspace_id, path, flow_version, deployment_msg) VALUES ($1, $2, $3, $4) ON CONFLICT (workspace_id, path, flow_version) WHERE flow_version IS NOT NULL DO UPDATE SET deployment_msg = $4", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Int8", + "Text" + ] + }, + "nullable": [] + }, + "hash": "afc7c23c057748f6d4a61dbef17e433b8875c6588b91e38e8141d12189118412" +} diff --git a/backend/.sqlx/query-e26e976a64d267557c528ce1bf006b97de9f89e53cb36a3b0ed11e7c18caf55f.json b/backend/.sqlx/query-bafff2205d9ca74b033d1d1faf6b1a3398ce387847d73f894dfe850e096326c0.json similarity index 52% rename from backend/.sqlx/query-e26e976a64d267557c528ce1bf006b97de9f89e53cb36a3b0ed11e7c18caf55f.json rename to backend/.sqlx/query-bafff2205d9ca74b033d1d1faf6b1a3398ce387847d73f894dfe850e096326c0.json index 28fd60c026..dbd29f1977 100644 --- a/backend/.sqlx/query-e26e976a64d267557c528ce1bf006b97de9f89e53cb36a3b0ed11e7c18caf55f.json +++ b/backend/.sqlx/query-bafff2205d9ca74b033d1d1faf6b1a3398ce387847d73f894dfe850e096326c0.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE flow SET workspace_id = $1 WHERE workspace_id = $2", + "query": "UPDATE flow_version SET workspace_id = $1 WHERE workspace_id = $2", "describe": { "columns": [], "parameters": { @@ -11,5 +11,5 @@ }, "nullable": [] }, - "hash": "e26e976a64d267557c528ce1bf006b97de9f89e53cb36a3b0ed11e7c18caf55f" + "hash": "bafff2205d9ca74b033d1d1faf6b1a3398ce387847d73f894dfe850e096326c0" } diff --git a/backend/.sqlx/query-d860cd90e583e4666beb37b283bd990dbcb91b781728b5b6253051decfda83d6.json b/backend/.sqlx/query-d860cd90e583e4666beb37b283bd990dbcb91b781728b5b6253051decfda83d6.json new file mode 100644 index 0000000000..149c649285 --- /dev/null +++ b/backend/.sqlx/query-d860cd90e583e4666beb37b283bd990dbcb91b781728b5b6253051decfda83d6.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO flow \n (workspace_id, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at) \n SELECT workspace_id, $1, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at\n FROM flow\n WHERE path = $2 AND workspace_id = $3", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "d860cd90e583e4666beb37b283bd990dbcb91b781728b5b6253051decfda83d6" +} diff --git a/backend/.sqlx/query-f095f413aef2ad008f4cb3d6c9517ce8abb1fa3b34e6a62864ec37ab9442de9d.json b/backend/.sqlx/query-f095f413aef2ad008f4cb3d6c9517ce8abb1fa3b34e6a62864ec37ab9442de9d.json new file mode 100644 index 0000000000..877e23350f --- /dev/null +++ b/backend/.sqlx/query-f095f413aef2ad008f4cb3d6c9517ce8abb1fa3b34e6a62864ec37ab9442de9d.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM flow WHERE path LIKE ('u/' || $1 || '/%') AND workspace_id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "f095f413aef2ad008f4cb3d6c9517ce8abb1fa3b34e6a62864ec37ab9442de9d" +} diff --git a/backend/.sqlx/query-fcf885a2214d5ae47e652f6c003cf537e011efd44c256301357079283dfc8e02.json b/backend/.sqlx/query-fcf885a2214d5ae47e652f6c003cf537e011efd44c256301357079283dfc8e02.json new file mode 100644 index 0000000000..cdbd4e08d1 --- /dev/null +++ b/backend/.sqlx/query-fcf885a2214d5ae47e652f6c003cf537e011efd44c256301357079283dfc8e02.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE flow SET versions = array_append(versions, $1) WHERE path = $2 AND workspace_id = $3", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Int8", + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "fcf885a2214d5ae47e652f6c003cf537e011efd44c256301357079283dfc8e02" +} diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index ff6c486fd8..a6ea4d7c72 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -c5b7088087d1c3528380e57ab5820e7e14028da9 +c7ef7778d68c86977a5d803fe98c418ef6a485c7 diff --git a/backend/migrations/20240630102146_flow_versioning.down.sql b/backend/migrations/20240630102146_flow_versioning.down.sql new file mode 100644 index 0000000000..f6b2c1a45b --- /dev/null +++ b/backend/migrations/20240630102146_flow_versioning.down.sql @@ -0,0 +1,7 @@ +-- Add down migration script here +DROP INDEX deployment_metadata_flow; +ALTER TABLE deployment_metadata DROP COLUMN flow_version; +create index if not exists deployment_metadata_flow on deployment_metadata (workspace_id, path) where script_hash is null and app_version is null; + +alter table flow drop column versions; +DROP TABLE flow_version; \ No newline at end of file diff --git a/backend/migrations/20240630102146_flow_versioning.up.sql b/backend/migrations/20240630102146_flow_versioning.up.sql new file mode 100644 index 0000000000..fdda7b8f99 --- /dev/null +++ b/backend/migrations/20240630102146_flow_versioning.up.sql @@ -0,0 +1,55 @@ +-- Add up migration script here + +-- create flow_version table with index +CREATE TABLE flow_version ( + id bigserial PRIMARY KEY, + workspace_id varchar(50) NOT NULL, + path varchar(255) NOT NULL, + value jsonb, + schema json, + created_by varchar(50) NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + FOREIGN KEY (workspace_id, path) REFERENCES flow (workspace_id, path) ON DELETE CASCADE +); +CREATE INDEX index_flow_version_path_created_at ON flow_version (path, created_at); + + +-- add versions column to flow +ALTER TABLE flow ADD COLUMN versions bigint[] NOT NULL DEFAULT '{}'::bigint[]; +-- create flow_version records for existing flows and update flow versions +INSERT INTO flow_version (workspace_id, path, value, schema, created_by, created_at) SELECT workspace_id, path, value, schema, edited_by, edited_at FROM flow; +UPDATE flow +SET versions = subquery.versions +FROM ( + SELECT + path, + workspace_id, + array_agg(id) AS versions + FROM + flow_version + GROUP BY + path, + workspace_id +) subquery +WHERE + flow.path = subquery.path + AND flow.workspace_id = subquery.workspace_id; + +-- add flow_version column to deployment_metadata +ALTER TABLE deployment_metadata ADD COLUMN flow_version int8; +-- populate flow_version column in deployment_metadata +UPDATE deployment_metadata +SET flow_version = fv.id +FROM flow_version fv +WHERE deployment_metadata.workspace_id = fv.workspace_id +AND deployment_metadata.path = fv.path +AND deployment_metadata.app_version IS NULL AND deployment_metadata.script_hash IS NULL; +-- update flow metadata index to include flow_verison +DROP INDEX IF EXISTS deployment_metadata_flow; +CREATE UNIQUE INDEX IF NOT EXISTS deployment_metadata_flow ON deployment_metadata (workspace_id, path, flow_version) WHERE flow_version IS NOT NULL; + +-- make sure the windmill_user and windmill_admin roles have access to the new tables +GRANT ALL ON flow_version TO windmill_user; +GRANT ALL ON flow_version_id_seq TO windmill_user; +GRANT ALL ON flow_version TO windmill_admin; +GRANT ALL ON flow_version_id_seq TO windmill_admin; \ No newline at end of file diff --git a/backend/tests/fixtures/schedule.sql b/backend/tests/fixtures/schedule.sql index 4df1690aa3..cfaa195c8b 100644 --- a/backend/tests/fixtures/schedule.sql +++ b/backend/tests/fixtures/schedule.sql @@ -41,12 +41,22 @@ export async function main() { '', 'f/system/schedule_recovery_handler', -28028598712388160, 'deno', ''); -INSERT INTO public.flow(workspace_id, edited_by, value, schema, summary, description, path) VALUES ( +INSERT INTO public.flow(workspace_id, summary, description, path, versions, schema, value, edited_by) VALUES ( 'test-workspace', -'system', -'{"modules": [{"id": "a", "value": {"path": "f/system/failing_script", "type": "script", "input_transforms": {"fail": {"expr": "flow_input.fail", "type": "javascript"}}}}]}', +'', +'', +'f/system/failing_flow', +'{1}', '{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{"fail":{"default":true,"description":"","type":"boolean","format":""}},"required":[],"type":"object"}', -'', -'', -'f/system/failing_flow' +'{"modules": [{"id": "a", "value": {"path": "f/system/failing_script", "type": "script", "input_transforms": {"fail": {"expr": "flow_input.fail", "type": "javascript"}}}}]}', +'system' +); + +INSERT INTO public.flow_version(id, workspace_id, path, schema, value, created_by) VALUES ( +1, +'test-workspace', +'f/system/failing_flow', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{"fail":{"default":true,"description":"","type":"boolean","format":""}},"required":[],"type":"object"}', +'{"modules": [{"id": "a", "value": {"path": "f/system/failing_script", "type": "script", "input_transforms": {"fail": {"expr": "flow_input.fail", "type": "javascript"}}}}]}', +'system' ); \ No newline at end of file diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index cdb7f2928f..29377a60bd 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -3546,6 +3546,13 @@ paths: in: query schema: type: boolean + - name: with_deployment_msg + description: | + (default false) + include deployment message + in: query + schema: + type: boolean responses: "200": description: All scripts @@ -4288,6 +4295,13 @@ paths: in: query schema: type: boolean + - name: with_deployment_msg + description: | + (default false) + include deployment message + in: query + schema: + type: boolean responses: "200": description: All flow @@ -4305,6 +4319,84 @@ paths: draft_only: type: boolean + /w/{workspace}/flows/history/p/{path}: + get: + summary: get flow history by path + operationId: getFlowHistory + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/ScriptPath" + tags: + - flow + responses: + "200": + description: Flow history + content: + application/json: + schema: + type: array + items: + $ref: "#/components/schemas/FlowVersion" + + /w/{workspace}/flows/get/v/{version}/p/{path}: + get: + summary: get flow version + operationId: getFlowVersion + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - type: string + name: version + in: path + required: true + schema: + type: number + - $ref: "#/components/parameters/ScriptPath" + tags: + - flow + responses: + "200": + description: flow details + content: + application/json: + schema: + $ref: "#/components/schemas/Flow" + + /w/{workspace}/flows/history_update/v/{version}/p/{path}: + post: + summary: update flow history + operationId: updateFlowHistory + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - type: string + name: version + in: path + required: true + schema: + type: number + - $ref: "#/components/parameters/ScriptPath" + requestBody: + description: Flow deployment message + required: true + content: + application/json: + schema: + type: object + properties: + deployment_msg: + type: string + required: + - deployment_msg + tags: + - flow + responses: + "200": + description: success + content: + text/plain: + schema: + type: string + + /w/{workspace}/flows/get/{path}: get: summary: get flow by path @@ -4648,6 +4740,13 @@ paths: in: query schema: type: boolean + - name: with_deployment_msg + description: | + (default false) + include deployment message + in: query + schema: + type: boolean responses: "200": description: All apps @@ -10615,6 +10714,20 @@ components: required: - version + FlowVersion: + type: object + properties: + id: + type: integer + created_at: + type: string + format: date-time + deployment_msg: + type: string + required: + - id + - created_at + SlackToken: type: object properties: diff --git a/backend/windmill-api/src/apps.rs b/backend/windmill-api/src/apps.rs index 675178a660..bf7ad1a730 100644 --- a/backend/windmill-api/src/apps.rs +++ b/backend/windmill-api/src/apps.rs @@ -91,6 +91,9 @@ pub struct ListableApp { pub has_draft: bool, #[serde(skip_serializing_if = "Option::is_none")] pub draft_only: Option, + #[sqlx(default)] + #[serde(skip_serializing_if = "Option::is_none")] + pub deployment_msg: Option, } #[derive(FromRow, Serialize, Deserialize)] @@ -291,6 +294,13 @@ async fn list_apps( sqlb.and_where("app.draft_only IS NOT TRUE"); } + if lq.with_deployment_msg.unwrap_or(false) { + sqlb.join("deployment_metadata dm") + .left() + .on("dm.app_version = app.versions[array_upper(app.versions, 1)]") + .fields(&["dm.deployment_msg"]); + } + let sql = sqlb.sql().map_err(|e| Error::InternalErr(e.to_string()))?; let mut tx = user_db.begin(&authed).await?; let rows = sqlx::query_as::<_, ListableApp>(&sql) diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index ecf7faac1f..e5a27972f6 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -56,6 +56,12 @@ pub fn workspaced_service() -> Router { .route("/get/draft/*path", get(get_flow_by_path_w_draft)) .route("/exists/*path", get(exists_flow_by_path)) .route("/list_paths", get(list_paths)) + .route("/history/p/*path", get(get_flow_history)) + .route( + "/history_update/v/:version/p/*path", + post(update_flow_history), + ) + .route("/get/v/:version/p/*path", get(get_flow_version)) .route( "/toggle_workspace_error_handler/*path", post(toggle_workspace_error_handler), @@ -86,7 +92,10 @@ async fn list_search_flows( let mut tx = user_db.begin(&authed).await?; let rows = sqlx::query_as::<_, SearchFlow>( - "SELECT path, value from flow WHERE workspace_id = $1 LIMIT $2", + "SELECT flow.path, flow_version.value + FROM flow + LEFT JOIN flow_version ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE workspace_id = $1 LIMIT $2", ) .bind(&w_id) .bind(n) @@ -113,8 +122,8 @@ async fn list_flows( "o.path", "summary", "description", - "edited_by", - "edited_at", + "fv.created_by as edited_by", + "fv.created_at as edited_at", "archived", "extra_perms", "favorite.path IS NOT NULL as starred", @@ -133,8 +142,13 @@ async fn list_flows( .on( "draft.path = o.path AND draft.workspace_id = o.workspace_id AND draft.typ = 'flow'" ) + .left() + .join("flow_version fv") + .on( + "fv.id = o.versions[array_upper(o.versions, 1)]" + ) .order_desc("favorite.path IS NOT NULL") - .order_by("edited_at", lq.order_desc.unwrap_or(true)) + .order_by("fv.created_at", lq.order_desc.unwrap_or(true)) .and_where("o.workspace_id = ?".bind(&w_id)) .offset(offset) .limit(per_page) @@ -149,7 +163,7 @@ async fn list_flows( sqlb.and_where_eq("o.path", "?".bind(p)); } if let Some(cb) = &lq.edited_by { - sqlb.and_where_eq("edited_by", "?".bind(cb)); + sqlb.and_where_eq("fv.created_by", "?".bind(cb)); } if lq.starred_only.unwrap_or(false) { sqlb.and_where_is_not_null("favorite.path"); @@ -159,6 +173,13 @@ async fn list_flows( sqlb.and_where("o.draft_only IS NOT TRUE"); } + if lq.with_deployment_msg.unwrap_or(false) { + sqlb.join("deployment_metadata dm") + .left() + .on("dm.flow_version = o.versions[array_upper(o.versions, 1)]") + .fields(&["dm.deployment_msg"]); + } + let sql = sqlb.sql().map_err(|e| Error::InternalErr(e.to_string()))?; let mut tx = user_db.begin(&authed).await?; let rows = sqlx::query_as::<_, ListableFlow>(&sql) @@ -324,24 +345,47 @@ async fn create_flow( check_path_conflict(tx.transaction_mut(), &w_id, &nf.path).await?; check_schedule_conflict(tx.transaction_mut(), &w_id, &nf.path).await?; + let schema_str = nf.schema.and_then(|x| serde_json::to_string(&x.0).ok()); + sqlx::query!( - "INSERT INTO flow (workspace_id, path, summary, description, value, edited_by, edited_at, \ - schema, dependency_job, draft_only, tag, dedicated_worker, visible_to_runner_only) VALUES ($1, $2, $3, $4, $5, $6, now(), $7::text::json, NULL, $8, $9, $10, $11)", + "INSERT INTO flow (workspace_id, path, summary, description, \ + dependency_job, draft_only, tag, dedicated_worker, visible_to_runner_only, value, schema, edited_by, edited_at) + VALUES ($1, $2, $3, $4, NULL, $5, $6, $7, $8, $9, $10::text::json, $11, now())", w_id, nf.path, nf.summary, nf.description.unwrap_or_else(String::new), - nf.value, - &authed.username, - nf.schema.and_then(|x| serde_json::to_string(&x.0).ok()), nf.draft_only, nf.tag, nf.dedicated_worker, nf.visible_to_runner_only.unwrap_or(false), + nf.value, + schema_str, + &authed.username, ) .execute(&mut tx) .await?; + let version = sqlx::query_scalar!( + "INSERT INTO flow_version (workspace_id, path, value, schema, created_by) + VALUES ($1, $2, $3, $4::text::json, $5) + RETURNING id", + w_id, + nf.path, + nf.value, + schema_str, + &authed.username, + ) + .fetch_one(&mut tx) + .await?; + + sqlx::query!( + "UPDATE flow SET versions = array_append(versions, $1) WHERE path = $2 AND workspace_id = $3", + version, + nf.path, + w_id + ).execute(&mut tx).await?; + sqlx::query!( "DELETE FROM draft WHERE path = $1 AND workspace_id = $2 AND typ = 'flow'", nf.path, @@ -379,6 +423,7 @@ async fn create_flow( JobPayload::FlowDependencies { path: nf.path.clone(), dedicated_worker: nf.dedicated_worker, + version: version, }, args.into(), &authed.username, @@ -454,6 +499,110 @@ pub async fn require_is_writer(authed: &ApiAuthed, path: &str, w_id: &str, db: D .await; } +#[derive(Serialize)] +pub struct FlowVersion { + pub id: i64, + pub created_at: chrono::DateTime, + #[serde(skip_serializing_if = "Option::is_none")] + pub deployment_msg: Option, +} + +async fn get_flow_history( + authed: ApiAuthed, + Extension(user_db): Extension, + Path((w_id, path)): Path<(String, StripPath)>, +) -> JsonResult> { + let path = path.to_path(); + let mut tx = user_db.begin(&authed).await?; + + let flows = sqlx::query_as!( + FlowVersion, + "SELECT flow_version.id, flow_version.created_at, deployment_metadata.deployment_msg FROM flow_version + LEFT JOIN deployment_metadata ON flow_version.id = deployment_metadata.flow_version + WHERE flow_version.path = $1 AND flow_version.workspace_id = $2 + ORDER BY flow_version.created_at DESC", + path, + w_id + ) + .fetch_all(&mut *tx) + .await?; + tx.commit().await?; + + Ok(Json(flows)) +} + +async fn get_flow_version( + authed: ApiAuthed, + Extension(user_db): Extension, + Path((w_id, version, path)): Path<(String, i64, StripPath)>, +) -> JsonResult { + let path = path.to_path(); + let mut tx = user_db.begin(&authed).await?; + + let flow = sqlx::query_as::<_, Flow>( + "SELECT flow.workspace_id, flow.path, flow.summary, flow.description, flow.archived, flow.extra_perms, flow.draft_only, flow.dedicated_worker, flow.tag, flow.ws_error_handler_muted, flow.timeout, flow.visible_to_runner_only, flow_version.schema, flow_version.value, flow_version.created_at as edited_at, flow_version.created_by as edited_by + FROM flow + LEFT JOIN flow_version ON flow_version.path = flow.path AND flow_version.workspace_id = flow.workspace_id + WHERE flow.path = $1 AND flow.workspace_id = $2 AND flow_version.id = $3", + ) + .bind(path) + .bind(w_id) + .bind(version) + .fetch_optional(&mut *tx) + .await?; + + tx.commit().await?; + + let flow = not_found_if_none(flow, "Flow version", version.to_string())?; + + Ok(Json(flow)) +} + +#[derive(Deserialize)] +pub struct FlowHistoryUpdate { + pub deployment_msg: String, +} + +async fn update_flow_history( + authed: ApiAuthed, + Extension(user_db): Extension, + Path((w_id, version, path)): Path<(String, i64, StripPath)>, + Json(history_update): Json, +) -> Result<()> { + let path = path.to_path(); + let mut tx = user_db.begin(&authed).await?; + let path_o = sqlx::query_scalar!( + "SELECT flow.path FROM flow + LEFT JOIN flow_version + ON flow_version.path = flow.path AND flow_version.workspace_id = flow.workspace_id + WHERE flow.path = $1 AND flow.workspace_id = $2 AND flow_version.id = $3", + path, + w_id, + version + ) + .fetch_optional(&mut *tx) + .await?; + + if path_o.is_none() { + tx.commit().await?; + return Err(Error::NotFound( + format!("Flow version {version} for path {path} not found").to_string(), + )); + } + + sqlx::query!( + "INSERT INTO deployment_metadata (workspace_id, path, flow_version, deployment_msg) VALUES ($1, $2, $3, $4) ON CONFLICT (workspace_id, path, flow_version) WHERE flow_version IS NOT NULL DO UPDATE SET deployment_msg = $4", + w_id, + path_o.unwrap(), + version, + history_update.deployment_msg, + ) + .fetch_optional(&mut *tx) + .await?; + tx.commit().await?; + return Ok(()); +} + async fn update_flow( authed: ApiAuthed, Extension(user_db): Extension, @@ -492,26 +641,99 @@ async fn update_flow( .fetch_optional(&mut tx) .await?; let old_dep_job = not_found_if_none(old_dep_job, "Flow", flow_path)?; + + let is_new_path = nf.path != flow_path; + + let schema_str = schema.and_then(|x| serde_json::to_string(&x).ok()); + sqlx::query!( - "UPDATE flow SET path = $1, summary = $2, description = $3, value = $4, edited_by = $5, \ - edited_at = now(), schema = $6::text::json, dependency_job = NULL, draft_only = NULL, tag = $9, dedicated_worker = $10, visible_to_runner_only = $11 - WHERE path = $7 AND workspace_id = $8", - nf.path, + "UPDATE flow SET path = $1, summary = $2, description = $3,\ + dependency_job = NULL, draft_only = NULL, tag = $4, dedicated_worker = $5, visible_to_runner_only = $6, \ + value = $7, schema = $8::text::json, edited_by = $9, edited_at = now() + WHERE path = $10 AND workspace_id = $11", + if is_new_path { flow_path } else { &nf.path }, // if new path, do not rename directly (to avoid flow_version foreign key constraint) nf.summary, nf.description.unwrap_or_else(String::new), - nf.value, - &authed.username, - schema.and_then(|x| serde_json::to_string(&x).ok()), - flow_path, - w_id, nf.tag, nf.dedicated_worker, nf.visible_to_runner_only.unwrap_or(false), + nf.value, + schema_str, + &authed.username, + flow_path, + w_id, ) .execute(&mut tx) .await.map_err(|e| error::Error::InternalErr(format!("Error updating flow due to flow update: {e:#}")))?; - if nf.path != flow_path { + if is_new_path { + // if new path, must clone flow to new path and delete old flow for flow_version foreign key constraint + sqlx::query!( + "INSERT INTO flow + (workspace_id, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at) + SELECT workspace_id, $1, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at + FROM flow + WHERE path = $2 AND workspace_id = $3", + nf.path, + flow_path, + w_id + ) + .execute(&mut tx) + .await + .map_err(|e| { + error::Error::InternalErr(format!("Error updating flow due to create new flow: {e:#}")) + })?; + + sqlx::query!( + "UPDATE flow_version SET path = $1 WHERE path = $2 AND workspace_id = $3", + nf.path, + flow_path, + w_id + ) + .execute(&mut tx) + .await + .map_err(|e| { + error::Error::InternalErr(format!( + "Error updating flow due to updating flow history path: {e:#}" + )) + })?; + + sqlx::query!( + "DELETE FROM flow WHERE path = $1 AND workspace_id = $2", + flow_path, + w_id + ) + .execute(&mut tx) + .await + .map_err(|e| { + error::Error::InternalErr(format!( + "Error updating flow due to deleting old flow: {e:#}" + )) + })?; + } + + let version = sqlx::query_scalar!( + "INSERT INTO flow_version (workspace_id, path, value, schema, created_by) VALUES ($1, $2, $3, $4::text::json, $5) RETURNING id", + w_id, + nf.path, + nf.value, + schema_str, + &authed.username, + ) + .fetch_one(&mut tx) + .await + .map_err(|e| { + error::Error::InternalErr(format!( + "Error updating flow due to flow history insert: {e:#}" + )) + })?; + + sqlx::query!( + "UPDATE flow SET versions = array_append(versions, $1) WHERE path = $2 AND workspace_id = $3", + version, flow_path, w_id + ).execute(&mut tx).await?; + + if is_new_path { check_schedule_conflict(tx.transaction_mut(), &w_id, &nf.path).await?; if !authed.is_admin { @@ -596,6 +818,7 @@ async fn update_flow( JobPayload::FlowDependencies { path: nf.path.clone(), dedicated_worker: nf.dedicated_worker, + version: version, }, windmill_queue::PushArgs { args, extra: HashMap::new() }, &authed.username, @@ -658,7 +881,12 @@ async fn get_flow_by_path( let mut tx = user_db.begin(&authed).await?; let flow_o = - sqlx::query_as::<_, Flow>("SELECT * FROM flow WHERE path = $1 AND workspace_id = $2") + sqlx::query_as::<_, Flow>( + "SELECT flow.workspace_id, flow.path, flow.summary, flow.description, flow.archived, flow.extra_perms, flow.draft_only, flow.dedicated_worker, flow.tag, flow.ws_error_handler_muted, flow.timeout, flow.visible_to_runner_only, flow_version.schema, flow_version.value, flow_version.created_at as edited_at, flow_version.created_by as edited_by + FROM flow + LEFT JOIN flow_version ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.path = $1 AND flow.workspace_id = $2" + ) .bind(path) .bind(w_id) .fetch_optional(&mut *tx) @@ -700,10 +928,12 @@ async fn get_flow_by_path_w_draft( let mut tx = user_db.begin(&authed).await?; let flow_o = sqlx::query_as::<_, FlowWDraft>( - "SELECT flow.path, flow.summary, flow,description, flow.schema, flow.value, flow.extra_perms, flow.draft_only, flow.ws_error_handler_muted, flow.dedicated_worker, draft.value as draft, flow.tag, flow.visible_to_runner_only + "SELECT flow.path, flow.summary, flow,description, flow_version.schema, flow_version.value, flow.extra_perms, flow.draft_only, flow.ws_error_handler_muted, flow.dedicated_worker, draft.value as draft, flow.tag, flow.visible_to_runner_only FROM flow - LEFT JOIN draft ON - flow.path = draft.path AND draft.workspace_id = $2 AND draft.typ = 'flow' + LEFT JOIN draft + ON flow.path = draft.path AND draft.workspace_id = $2 AND draft.typ = 'flow' + LEFT JOIN flow_version + ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] WHERE flow.path = $1 AND flow.workspace_id = $2", ) .bind(path) @@ -777,7 +1007,11 @@ async fn archive_flow_by_path( &authed.username, &db, &w_id, - DeployedObject::Flow { path: path.to_string(), parent_path: Some(path.to_string()) }, + DeployedObject::Flow { + path: path.to_string(), + parent_path: Some(path.to_string()), + version: 0, // dummy version as it will not get inserted in db + }, Some(format!( "Flow '{}' {}", path, @@ -844,7 +1078,11 @@ async fn delete_flow_by_path( &authed.username, &db, &w_id, - DeployedObject::Flow { path: path.to_string(), parent_path: Some(path.to_string()) }, + DeployedObject::Flow { + path: path.to_string(), + parent_path: Some(path.to_string()), + version: 0, // dummy version as it will not get inserted in db + }, Some(format!("Flow '{}' deleted", path)), rsmq, true, diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 58e88b2772..b0bf34c104 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -3470,7 +3470,11 @@ async fn run_wait_result_flow_by_path_internal( let scheduled_for = run_query.get_scheduled_for(&db).await?; let (tag, dedicated_worker, early_return) = sqlx::query!( - "SELECT tag, dedicated_worker, value->>'early_return' as early_return from flow WHERE path = $1 and workspace_id = $2", + "SELECT tag, dedicated_worker, flow_version.value->>'early_return' as early_return + FROM flow + LEFT JOIN flow_version + ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.path = $1 and flow.workspace_id = $2", flow_path, w_id ) diff --git a/backend/windmill-api/src/scripts.rs b/backend/windmill-api/src/scripts.rs index 8656ace619..967df3f201 100644 --- a/backend/windmill-api/src/scripts.rs +++ b/backend/windmill-api/src/scripts.rs @@ -280,6 +280,13 @@ async fn list_scripts( sqlb.and_where_is_not_null("favorite.path"); } + if lq.with_deployment_msg.unwrap_or(false) { + sqlb.join("deployment_metadata dm") + .left() + .on("dm.script_hash = o.hash") + .fields(&["dm.deployment_msg"]); + } + let sql = sqlb.sql().map_err(|e| Error::InternalErr(e.to_string()))?; let mut tx = user_db.begin(&authed).await?; let rows = sqlx::query_as::<_, ListableScript>(&sql) diff --git a/backend/windmill-api/src/users.rs b/backend/windmill-api/src/users.rs index 2bcdf01983..d60bafec09 100644 --- a/backend/windmill-api/src/users.rs +++ b/backend/windmill-api/src/users.rs @@ -2476,7 +2476,11 @@ async fn get_all_runnables( })?; let mut tx = db.clone().begin(&nauthed).await?; let flows = sqlx::query!( - "SELECT workspace_id as workspace, path, summary, description, schema FROM flow WHERE workspace_id = $1", workspace + "SELECT flow.workspace_id as workspace, flow.path, summary, description, flow_version.schema + FROM flow + LEFT JOIN flow_version ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.workspace_id = $1", + workspace ) .fetch_all(&mut *tx) .await?; @@ -2916,7 +2920,19 @@ async fn update_username_in_workpsace<'c>( // ---- flows ---- sqlx::query!( - r#"UPDATE flow SET path = REGEXP_REPLACE(path,'u/' || $2 || '/(.*)','u/' || $1 || '/\1') WHERE path LIKE ('u/' || $2 || '/%') AND workspace_id = $3"#, + r#"INSERT INTO flow + (workspace_id, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at) + SELECT workspace_id, REGEXP_REPLACE(path,'u/' || $2 || '/(.*)','u/' || $1 || '/\1'), summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at + FROM flow + WHERE path LIKE ('u/' || $2 || '/%') AND workspace_id = $3"#, + new_username, + old_username, + w_id + ).execute(&mut **tx) + .await?; + + sqlx::query!( + r#"UPDATE flow_version SET path = REGEXP_REPLACE(path,'u/' || $2 || '/(.*)','u/' || $1 || '/\1') WHERE path LIKE ('u/' || $2 || '/%') AND workspace_id = $3"#, new_username, old_username, w_id @@ -2925,14 +2941,12 @@ async fn update_username_in_workpsace<'c>( .await?; sqlx::query!( - "UPDATE flow SET edited_by = $1 WHERE edited_by = $2 AND workspace_id = $3", - new_username, + "DELETE FROM flow WHERE path LIKE ('u/' || $1 || '/%') AND workspace_id = $2", old_username, w_id ) .execute(&mut **tx) - .await - .unwrap(); + .await?; sqlx::query!( "UPDATE flow SET extra_perms = extra_perms - ('u/' || $2) || jsonb_build_object(('u/' || $1), extra_perms->('u/' || $2)) WHERE extra_perms ? ('u/' || $2) AND workspace_id = $3", diff --git a/backend/windmill-api/src/workspaces.rs b/backend/windmill-api/src/workspaces.rs index 0dff176ff5..51e4d1df39 100644 --- a/backend/windmill-api/src/workspaces.rs +++ b/backend/windmill-api/src/workspaces.rs @@ -2572,7 +2572,10 @@ async fn tarball_workspace( { let flows = sqlx::query_as::<_, Flow>( - "SELECT * FROM flow WHERE workspace_id = $1 AND archived = false", + "SELECT flow.workspace_id, flow.path, flow.summary, flow.description, flow.archived, flow.extra_perms, flow.draft_only, flow.dedicated_worker, flow.tag, flow.ws_error_handler_muted, flow.timeout, flow.visible_to_runner_only, flow_version.schema, flow_version.value, flow_version.created_at as edited_at, flow_version.created_by as edited_by + FROM flow + LEFT JOIN flow_version ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.workspace_id = $1 AND flow.archived = false", ) .bind(&w_id) .fetch_all(&mut *tx) @@ -2970,7 +2973,10 @@ async fn change_workspace_id( .await?; sqlx::query!( - "UPDATE flow SET workspace_id = $1 WHERE workspace_id = $2", + "INSERT INTO flow + (workspace_id, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at) + SELECT $1, path, summary, description, archived, extra_perms, dependency_job, draft_only, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, concurrency_key, versions, value, schema, edited_by, edited_at + FROM flow WHERE workspace_id = $2", &rw.new_id, &old_id ) @@ -2978,13 +2984,17 @@ async fn change_workspace_id( .await?; sqlx::query!( - "UPDATE folder SET workspace_id = $1 WHERE workspace_id = $2", + "UPDATE flow_version SET workspace_id = $1 WHERE workspace_id = $2", &rw.new_id, &old_id ) .execute(&mut *tx) .await?; + sqlx::query!("DELETE FROM flow WHERE workspace_id = $1", &old_id) + .execute(&mut *tx) + .await?; + // have to duplicate group_ with new workspace id because of foreign key constraint sqlx::query!( "INSERT INTO group_ SELECT $1, name, summary, extra_perms FROM group_ WHERE workspace_id = $2", @@ -3007,6 +3017,14 @@ async fn change_workspace_id( .execute(&mut *tx) .await?; + sqlx::query!( + "UPDATE folder SET workspace_id = $1 WHERE workspace_id = $2", + &rw.new_id, + &old_id + ) + .execute(&mut *tx) + .await?; + sqlx::query!( "UPDATE input SET workspace_id = $1 WHERE workspace_id = $2", &rw.new_id, diff --git a/backend/windmill-common/src/apps.rs b/backend/windmill-common/src/apps.rs index fb7499203b..acbf720f57 100644 --- a/backend/windmill-common/src/apps.rs +++ b/backend/windmill-common/src/apps.rs @@ -14,4 +14,5 @@ pub struct ListAppQuery { pub path_exact: Option, pub path_start: Option, pub include_draft_only: Option, + pub with_deployment_msg: Option, } diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index 37473b333a..575f44c3a5 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -20,7 +20,7 @@ use crate::{ scripts::{Schema, ScriptHash, ScriptLang}, }; -#[derive(Serialize, sqlx::FromRow)] +#[derive(Serialize, Deserialize, sqlx::FromRow)] pub struct Flow { pub workspace_id: String, pub path: String, @@ -62,6 +62,9 @@ pub struct ListableFlow { pub draft_only: Option, #[serde(skip_serializing_if = "Option::is_none")] pub ws_error_handler_muted: Option, + #[sqlx(default)] + #[serde(skip_serializing_if = "Option::is_none")] + pub deployment_msg: Option, } #[derive(Deserialize, sqlx::FromRow)] @@ -73,7 +76,6 @@ pub struct NewFlow { pub schema: Option, pub draft_only: Option, pub tag: Option, - pub ws_error_handler_muted: Option, pub dedicated_worker: Option, pub timeout: Option, pub deployment_message: Option, @@ -552,6 +554,7 @@ pub struct ListFlowQuery { pub order_desc: Option, pub starred_only: Option, pub include_draft_only: Option, + pub with_deployment_msg: Option, } pub fn add_virtual_items_if_necessary(modules: &mut Vec) { diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index b29a04633a..300d204415 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -318,6 +318,7 @@ pub enum JobPayload { FlowDependencies { path: String, dedicated_worker: Option, + version: i64, }, AppDependencies { path: String, diff --git a/backend/windmill-common/src/scripts.rs b/backend/windmill-common/src/scripts.rs index bd1f2b0abc..a88d357667 100644 --- a/backend/windmill-common/src/scripts.rs +++ b/backend/windmill-common/src/scripts.rs @@ -206,6 +206,9 @@ pub struct ListableScript { pub no_main_func: Option, #[serde(skip_serializing_if = "is_false")] pub use_codebase: bool, + #[sqlx(default)] + #[serde(skip_serializing_if = "Option::is_none")] + pub deployment_msg: Option, } fn is_false(x: &bool) -> bool { @@ -338,6 +341,7 @@ pub struct ListScriptQuery { pub starred_only: Option, pub include_without_main: Option, pub include_draft_only: Option, + pub with_deployment_msg: Option, } pub fn to_i64(s: &str) -> crate::error::Result { diff --git a/backend/windmill-git-sync/src/lib.rs b/backend/windmill-git-sync/src/lib.rs index 9bd5b0017d..707203accf 100644 --- a/backend/windmill-git-sync/src/lib.rs +++ b/backend/windmill-git-sync/src/lib.rs @@ -18,7 +18,7 @@ pub type DB = Pool; #[derive(Clone, Debug)] pub enum DeployedObject { Script { hash: ScriptHash, path: String, parent_path: Option }, - Flow { path: String, parent_path: Option }, + Flow { path: String, parent_path: Option, version: i64 }, App { path: String, version: i64, parent_path: Option }, Folder { path: String }, Resource { path: String, parent_path: Option }, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 730f251779..323d6d37ad 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -2080,7 +2080,11 @@ pub async fn custom_concurrency_key( async fn legacy_concurrency_key(db: &Pool, queued_job: &QueuedJob) -> Option { let r = if queued_job.is_flow() { sqlx::query_scalar!( - "SELECT value->>'concurrency_key' FROM flow WHERE path = $1 AND workspace_id = $2", + "SELECT flow_version.value->>'concurrency_key' + FROM flow + LEFT JOIN flow_version + ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.path = $1 AND flow.workspace_id = $2", queued_job.script_path, queued_job.workspace_id ) @@ -3270,13 +3274,10 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( None, None, ), - JobPayload::FlowDependencies { path, dedicated_worker } => { + JobPayload::FlowDependencies { path, dedicated_worker, version } => { let value_json = fetch_scalar_isolated!( - sqlx::query_as::<_, FlowRawValue>( - "SELECT value FROM flow WHERE path = $1 AND workspace_id = $2", - ) - .bind(&path) - .bind(&workspace_id), + sqlx::query_as::<_, FlowRawValue>("SELECT value FROM flow_version WHERE id = $1",) + .bind(&version), tx )? .ok_or_else(|| Error::InternalErr(format!("not found flow at path {:?}", path)))?; @@ -3287,7 +3288,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( )) })?; ( - None, + Some(version), Some(path), None, JobKind::FlowDependencies, @@ -3441,7 +3442,10 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( JobPayload::Flow { path, dedicated_worker } => { let value_json = fetch_scalar_isolated!( sqlx::query_as::<_, FlowRawValue>( - "SELECT value FROM flow WHERE path = $1 AND workspace_id = $2", + "SELECT flow_version.value FROM flow + LEFT JOIN flow_version + ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.path = $1 AND flow.workspace_id = $2", ) .bind(&path) .bind(&workspace_id), diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index f31102ced8..1dc06a156d 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -400,11 +400,15 @@ pub async fn create_dedicated_worker_map( if let Some(flow_path) = _wp.path.strip_prefix("flow/") { is_flow_worker = true; let value = sqlx::query_scalar!( - "SELECT value FROM flow WHERE path = $1 AND workspace_id = $2", + "SELECT flow_version.value + FROM flow + LEFT JOIN flow_version + ON flow_version.id = flow.versions[array_upper(flow.versions, 1)] + WHERE flow.path = $1 AND flow.workspace_id = $2", flow_path, _wp.workspace_id ) - .fetch_optional(db) + .fetch_one(db) .await; if let Ok(v) = value { if let Some(v) = v { diff --git a/backend/windmill-worker/src/worker_lockfiles.rs b/backend/windmill-worker/src/worker_lockfiles.rs index c789174bc0..c0ea748464 100644 --- a/backend/windmill-worker/src/worker_lockfiles.rs +++ b/backend/windmill-worker/src/worker_lockfiles.rs @@ -409,7 +409,34 @@ async fn trigger_dependents_to_recompute_dependencies< "nodes_to_relock".to_string(), to_raw_value(&s.importer_node_ids), ); - JobPayload::FlowDependencies { path: s.importer_path.clone(), dedicated_worker: None } + let r = sqlx::query_scalar!( + "SELECT versions[array_upper(versions, 1)] FROM flow WHERE path = $1 AND workspace_id = $2", + s.importer_path, + w_id, + ).fetch_one(db) + .await + .map_err(to_anyhow); + match r { + Ok(Some(version)) => JobPayload::FlowDependencies { + path: s.importer_path.clone(), + dedicated_worker: None, + version: version, + }, + Ok(None) => { + tracing::error!( + "no flow version found for path {path}", + path = s.importer_path + ); + continue; + } + Err(err) => { + tracing::error!( + "error getting latest deployed flow version for path {path}: {err}", + path = s.importer_path, + ); + continue; + } + } } else { tracing::error!( "unexpected importer kind: {kind} for path {path}", @@ -470,6 +497,15 @@ pub async fn handle_flow_dependency_job error::Result<()> { let job_path = job.script_path.clone().ok_or_else(|| { error::Error::InternalErr( - "Cannot resolve flow dependencies for flow without path".to_string(), + "Cannot resolve app dependencies for app without path".to_string(), ) })?; let id = job .script_hash .clone() - .ok_or_else(|| Error::InternalErr("Flow Dependency requires script hash".to_owned()))? + .ok_or_else(|| Error::InternalErr("App Dependency requires script hash".to_owned()))? .0; let value = sqlx::query_scalar!("SELECT value FROM app_version WHERE id = $1", id) .fetch_optional(db) @@ -1045,7 +1088,7 @@ pub async fn handle_app_dependency_job @@ -139,6 +146,7 @@ }) }} disabled={email === undefined || !username} + {loading} > Confirm username change diff --git a/frontend/src/lib/components/CustomPopover.svelte b/frontend/src/lib/components/CustomPopover.svelte new file mode 100644 index 0000000000..8392ca4c0c --- /dev/null +++ b/frontend/src/lib/components/CustomPopover.svelte @@ -0,0 +1,81 @@ + + +{#if notClickable} + + + + +{:else} + +{/if} +{#if showTooltip && !disablePopup} + + +
+ +
+
+{/if} diff --git a/frontend/src/lib/components/FlowBuilder.svelte b/frontend/src/lib/components/FlowBuilder.svelte index da6a52d6ae..5a9aa8eca5 100644 --- a/frontend/src/lib/components/FlowBuilder.svelte +++ b/frontend/src/lib/components/FlowBuilder.svelte @@ -27,9 +27,9 @@ sleep } from '$lib/utils' import { sendUserToast } from '$lib/toast' - import type { Drawer } from '$lib/components/common' + import { Drawer } from '$lib/components/common' - import { setContext, tick } from 'svelte' + import { setContext, tick, type ComponentType } from 'svelte' import { writable, type Writable } from 'svelte/store' import CenteredPage from './CenteredPage.svelte' import { Badge, Button, UndoRedo } from './common' @@ -43,7 +43,17 @@ import { loadFlowSchedule, type Schedule } from './flows/scheduleUtils' import type { FlowEditorContext, FlowInput } from './flows/types' import { cleanInputs, emptyFlowModuleState } from './flows/utils' - import { Calendar, Pen, Save, DiffIcon } from 'lucide-svelte' + import { + Calendar, + Pen, + Save, + DiffIcon, + MoreVertical, + HistoryIcon, + FileJson, + type Icon, + CornerDownLeft + } from 'lucide-svelte' import { createEventDispatcher } from 'svelte' import Awareness from './Awareness.svelte' import { getAllModules } from './flows/flowExplorer' @@ -64,6 +74,11 @@ import FlowTutorials from './FlowTutorials.svelte' import { ignoredTutorials } from './tutorials/ignoredTutorials' import type DiffDrawer from './DiffDrawer.svelte' + import FlowHistory from './flows/FlowHistory.svelte' + import ButtonDropdown from './common/button/ButtonDropdown.svelte' + import { MenuItem } from '@rgossiaux/svelte-headlessui' + import { twMerge } from 'tailwind-merge' + import CustomPopover from './CustomPopover.svelte' import Summary from './Summary.svelte' export let initialPath: string = '' @@ -208,7 +223,7 @@ ) } - async function saveFlow(): Promise { + async function saveFlow(deploymentMsg?: string): Promise { loadingSave = true try { const flow = cleanInputs($flowStore) @@ -234,7 +249,8 @@ ws_error_handler_muted: flow.ws_error_handler_muted, tag: flow.tag, dedicated_worker: flow.dedicated_worker, - visible_to_runner_only: flow.visible_to_runner_only + visible_to_runner_only: flow.visible_to_runner_only, + deployment_message: deploymentMsg || undefined } }) if (enabled) { @@ -297,7 +313,8 @@ tag: flow.tag, dedicated_worker: flow.dedicated_worker, ws_error_handler_muted: flow.ws_error_handler_muted, - visible_to_runner_only: flow.visible_to_runner_only + visible_to_runner_only: flow.visible_to_runner_only, + deployment_message: deploymentMsg || undefined } }) } @@ -979,6 +996,9 @@ let renderCount = 0 let flowTutorials: FlowTutorials | undefined = undefined + let jsonViewerDrawer: Drawer | undefined = undefined + let flowHistory: FlowHistory | undefined = undefined + export function triggerTutorial() { const urlParams = new URLSearchParams(window.location.search) const tutorial = urlParams.get('tutorial') @@ -989,6 +1009,30 @@ flowTutorials?.runTutorialById('action') } } + + const moreItems: { + displayName: string + icon: ComponentType + action: () => void + disabled?: boolean + }[] = [ + { + displayName: 'Deployment History', + icon: HistoryIcon, + action: () => { + flowHistory?.open() + }, + disabled: newFlow + }, + { + displayName: 'Export', + icon: FileJson, + action: () => jsonViewerDrawer?.openDrawer() + } + ] + + let deploymentMsg = '' + let msgInput: HTMLInputElement | undefined = undefined @@ -998,6 +1042,10 @@ {#key renderCount} {#if !$userStore?.operator} + {#if $pathStore} + + {/if} + { applyCopilotFlowInputs() @@ -1078,10 +1126,40 @@ /> -
+
{#if $enterpriseLicense && !newFlow} {/if} +
+ + + + + + {#each moreItems as item} + +
+ + {item.displayName} +
+
+ {/each} +
+
+
+ { renderCount += 1 @@ -1119,8 +1197,6 @@ {abortController} /> - - - + + + + +
+ { + if (e.key === 'Enter') { + saveFlow(deploymentMsg) + } + }} + bind:this={msgInput} + /> + +
+
+
diff --git a/frontend/src/lib/components/LightweightArgInput.svelte b/frontend/src/lib/components/LightweightArgInput.svelte index f26ec735a6..872233291a 100644 --- a/frontend/src/lib/components/LightweightArgInput.svelte +++ b/frontend/src/lib/components/LightweightArgInput.svelte @@ -21,7 +21,6 @@ import Password from './Password.svelte' import ToggleButton from './common/toggleButton-v2/ToggleButton.svelte' import ToggleButtonGroup from './common/toggleButton-v2/ToggleButtonGroup.svelte' - import S3FilePicker from './S3FilePicker.svelte' import FileUpload from './common/fileUpload/FileUpload.svelte' export let css: ComponentCustomCSS<'schemaformcomponent'> | undefined = undefined diff --git a/frontend/src/lib/components/ScriptBuilder.svelte b/frontend/src/lib/components/ScriptBuilder.svelte index 717f224d60..25bec77170 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -33,6 +33,7 @@ Calendar, CheckCircle, Code, + CornerDownLeft, DiffIcon, Pen, Plus, @@ -59,6 +60,7 @@ import { type ScriptSchedule, loadScriptSchedule, defaultScriptLanguages } from '$lib/scripts' import DefaultScripts from './DefaultScripts.svelte' import { createEventDispatcher } from 'svelte' + import CustomPopover from './CustomPopover.svelte' import Summary from './Summary.svelte' export let script: NewScript @@ -206,7 +208,7 @@ } } - async function editScript(stay: boolean): Promise { + async function editScript(stay: boolean, deploymentMsg?: string): Promise { loadingSave = true try { try { @@ -247,7 +249,8 @@ timeout: script.timeout, concurrency_key: emptyString(script.concurrency_key) ? undefined : script.concurrency_key, visible_to_runner_only: script.visible_to_runner_only, - no_main_func: script.no_main_func + no_main_func: script.no_main_func, + deployment_message: deploymentMsg || undefined } }) @@ -456,6 +459,9 @@ let dirtyPath = false let selectedTab: 'metadata' | 'runtime' | 'ui' | 'schedule' = 'metadata' + + let deploymentMsg = '' + let msgInput: HTMLInputElement | undefined = undefined @@ -1101,15 +1107,41 @@ > - + + + + +
+ { + if (e.key === 'Enter') { + editScript(false, deploymentMsg) + } + }} + /> + +
+
+
diff --git a/frontend/src/lib/components/apps/editor/AppDeploymentHistory.svelte b/frontend/src/lib/components/apps/editor/AppDeploymentHistory.svelte index 193da625b6..a712b69087 100644 --- a/frontend/src/lib/components/apps/editor/AppDeploymentHistory.svelte +++ b/frontend/src/lib/components/apps/editor/AppDeploymentHistory.svelte @@ -5,11 +5,10 @@ import { sendUserToast } from '$lib/toast' import DeploymentHistory from './DeploymentHistory.svelte' - let appPath: string | undefined = undefined + export let appPath: string | undefined = undefined let historyBrowserDrawerOpen = false - export function open(appPath: string) { - appPath = appPath + export function open() { historyBrowserDrawerOpen = true } diff --git a/frontend/src/lib/components/common/table/AppRow.svelte b/frontend/src/lib/components/common/table/AppRow.svelte index 285a52bacf..c1ac3c4910 100644 --- a/frontend/src/lib/components/common/table/AppRow.svelte +++ b/frontend/src/lib/components/common/table/AppRow.svelte @@ -3,7 +3,7 @@ import type MoveDrawer from '$lib/components/MoveDrawer.svelte' import SharedBadge from '$lib/components/SharedBadge.svelte' import type ShareModal from '$lib/components/ShareModal.svelte' - import { AppService, type AppWithLastVersion, DraftService, type ListableApp } from '$lib/gen' + import { AppService, DraftService, type ListableApp } from '$lib/gen' import { userStore, workspaceStore } from '$lib/stores' import { createEventDispatcher } from 'svelte' import Button from '../button/Button.svelte' @@ -44,25 +44,16 @@ const dispatch = createEventDispatcher() let appExport: AppJsonEditor - let appDeploymentHistory: AppDeploymentHistory + let appDeploymentHistory: AppDeploymentHistory | undefined = undefined async function loadAppJson() { appExport.open(app.path) } - - async function loadDeployements() { - const napp: AppWithLastVersion = (await AppService.getAppByPath({ - workspace: $workspaceStore!, - path: app.path - })) as unknown as AppWithLastVersion - - appDeploymentHistory.open(napp.path) - } {#if menuOpen} - + {/if} loadDeployements(), + action: () => appDeploymentHistory?.open(), hide: $userStore?.operator }, { diff --git a/frontend/src/lib/components/common/table/FlowRow.svelte b/frontend/src/lib/components/common/table/FlowRow.svelte index 95fe7eefd1..50a3782f97 100644 --- a/frontend/src/lib/components/common/table/FlowRow.svelte +++ b/frontend/src/lib/components/common/table/FlowRow.svelte @@ -26,8 +26,10 @@ Share, Archive, Clipboard, - Eye + Eye, + HistoryIcon } from 'lucide-svelte' + import FlowHistory from '$lib/components/flows/FlowHistory.svelte' export let flow: Flow & { has_draft?: boolean; draft_only?: boolean; canWrite: boolean } export let marked: string | undefined @@ -66,10 +68,12 @@ } } let scheduleEditor: ScheduleEditor + let flowHistory: FlowHistory {#if menuOpen} goto('/schedules')} bind:this={scheduleEditor} /> + {/if} { + flowHistory.open() + }, + hide: $userStore?.operator + }, { displayName: 'Schedule', icon: Calendar, diff --git a/frontend/src/lib/components/flows/FlowHistory.svelte b/frontend/src/lib/components/flows/FlowHistory.svelte new file mode 100644 index 0000000000..0b26e5a3b5 --- /dev/null +++ b/frontend/src/lib/components/flows/FlowHistory.svelte @@ -0,0 +1,214 @@ + + + + { + drawer?.closeDrawer() + }} + > + + + +
+ {#if !loading} + {#if versions.length > 0} +
+ {#each versions ?? [] as version} + +
{ + selectedVersion = version + }} + > + + {#if emptyString(version.deployment_msg)}Version {version.id}{:else}{version.deployment_msg}{/if} + +
+ {/each} +
+ {:else} +
No items
+ {/if} + {:else} + + {/if} +
+
+
+ +
+ {#if selectedVersion} + {#if selected} +
+ + {#if deploymentMsgUpdateMode} +
+ {}} + on:keydown|stopPropagation + on:keypress|stopPropagation={({ key }) => { + if (key === 'Enter') updateDeploymentMsg(selectedVersion?.id) + }} + /> + + +
+ {:else} + {#if selectedVersion.deployment_msg} + {selectedVersion.deployment_msg} + {:else} + Deployed {displayDate(selected.edited_at)} by {selected.edited_by} + {/if} + + {/if} +
+
+ +
+ +
+ {:else} + + {/if} + {:else} +
Select a deployment version to see its details
+ {/if} +
+
+
+
+
diff --git a/frontend/src/lib/components/flows/header/FlowImportExportMenu.svelte b/frontend/src/lib/components/flows/header/FlowImportExportMenu.svelte index 3a98700c54..7ba42ba4aa 100644 --- a/frontend/src/lib/components/flows/header/FlowImportExportMenu.svelte +++ b/frontend/src/lib/components/flows/header/FlowImportExportMenu.svelte @@ -3,29 +3,16 @@ import DrawerContent from '$lib/components/common/drawer/DrawerContent.svelte' import FlowViewer from '$lib/components/FlowViewer.svelte' import { getContext } from 'svelte' - import { Button } from '../../common' import type { FlowEditorContext } from '../types' import { cleanInputs } from '../utils' - import { FileJson } from 'lucide-svelte' const { flowStore } = getContext('FlowEditorContext') - let jsonViewerDrawer: Drawer + export let drawer: Drawer | undefined - - - - jsonViewerDrawer.toggleDrawer()}> + + drawer?.toggleDrawer()}> {#if $flowStore} {/if} diff --git a/frontend/src/routes/(root)/(logged)/apps/edit/[...path]/+page.svelte b/frontend/src/routes/(root)/(logged)/apps/edit/[...path]/+page.svelte index c73af4cbda..a229837dba 100644 --- a/frontend/src/routes/(root)/(logged)/apps/edit/[...path]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/apps/edit/[...path]/+page.svelte @@ -180,6 +180,19 @@ } let diffDrawer: DiffDrawer + + function onRestore(ev: any) { + sendUserToast('App restored from previous deployment') + app = ev.detail + const app_ = structuredClone(app!) + savedApp = { + summary: app_.summary, + value: app_.value as App, + path: app_.path, + policy: app_.policy + } + redraw++ + } @@ -188,11 +201,7 @@ {#if app}
{ - sendUserToast('App restored from previous deployment') - app = e.detail - redraw++ - }} + on:restore={onRestore} summary={app.summary} app={app.value} path={app.path} diff --git a/frontend/src/routes/(root)/(logged)/flows/edit/[...path]/+page.svelte b/frontend/src/routes/(root)/(logged)/flows/edit/[...path]/+page.svelte index 0caac34efb..c40ec211da 100644 --- a/frontend/src/routes/(root)/(logged)/flows/edit/[...path]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/flows/edit/[...path]/+page.svelte @@ -205,6 +205,9 @@ const { path, selectedId } = e.detail goto(`/flows/edit/${path}?selected=${selectedId}`) }} + on:historyRestore={() => { + loadFlow() + }} {flowStore} {flowStateStore} initialPath={$page.params.path} diff --git a/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte b/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte index b7f500a880..cec3e426f4 100644 --- a/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte @@ -26,7 +26,8 @@ Columns, Pen, Eye, - Calendar + Calendar, + HistoryIcon } from 'lucide-svelte' import DetailPageHeader from '$lib/components/details/DetailPageHeader.svelte' @@ -42,6 +43,7 @@ import FlowGraphViewerStep from '$lib/components/FlowGraphViewerStep.svelte' import { loadFlowSchedule, type Schedule } from '$lib/components/flows/scheduleUtils' import GfmMarkdown from '$lib/components/GfmMarkdown.svelte' + import FlowHistory from '$lib/components/flows/FlowHistory.svelte' let flow: Flow | undefined let can_write = false @@ -231,6 +233,11 @@ }) if (can_write) { + menuItems.push({ + label: 'Deployments', + onclick: () => flowHistory?.open(), + Icon: HistoryIcon + }) menuItems.push({ label: flow.archived ? 'Unarchive' : 'Archive', onclick: () => flow?.path && archiveFlow(), @@ -266,6 +273,8 @@ let detailSelected = 'saved_inputs' let triggerSelected: 'webhooks' | 'schedule' | 'cli' = 'webhooks' + + let flowHistory: FlowHistory | undefined = undefined +{#if flow} + +{/if} {#if flow}