feat: flow versioning (#4009)

* feat: flow versioning

* fix: sqlx

* fix: update schedule test for flow versioning

* fix: with_deployment_msg + UI nits

* fix: nit

* fix: improve down migration

* patch: keep latest flow version in flow table for backward compat

* fix: app deployments in list view

* chore: update ee ref

* fix: merge

* fix: tests
This commit is contained in:
HugoCasa
2024-07-03 15:53:45 +02:00
committed by GitHub
parent 1a403f8f39
commit e50f1752da
57 changed files with 1464 additions and 234 deletions
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
+1 -1
View File
@@ -1 +1 @@
c5b7088087d1c3528380e57ab5820e7e14028da9
c7ef7778d68c86977a5d803fe98c418ef6a485c7
@@ -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;
@@ -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;
+16 -6
View File
@@ -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'
);
+113
View File
@@ -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:
+10
View File
@@ -91,6 +91,9 @@ pub struct ListableApp {
pub has_draft: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub draft_only: Option<bool>,
#[sqlx(default)]
#[serde(skip_serializing_if = "Option::is_none")]
pub deployment_msg: Option<String>,
}
#[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)
+264 -26
View File
@@ -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<chrono::Utc>,
#[serde(skip_serializing_if = "Option::is_none")]
pub deployment_msg: Option<String>,
}
async fn get_flow_history(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Vec<FlowVersion>> {
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<UserDB>,
Path((w_id, version, path)): Path<(String, i64, StripPath)>,
) -> JsonResult<Flow> {
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<UserDB>,
Path((w_id, version, path)): Path<(String, i64, StripPath)>,
Json(history_update): Json<FlowHistoryUpdate>,
) -> 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<UserDB>,
@@ -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,
+5 -1
View File
@@ -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
)
+7
View File
@@ -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)
+20 -6
View File
@@ -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",
+21 -3
View File
@@ -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,
+1
View File
@@ -14,4 +14,5 @@ pub struct ListAppQuery {
pub path_exact: Option<String>,
pub path_start: Option<String>,
pub include_draft_only: Option<bool>,
pub with_deployment_msg: Option<bool>,
}
+5 -2
View File
@@ -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<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ws_error_handler_muted: Option<bool>,
#[sqlx(default)]
#[serde(skip_serializing_if = "Option::is_none")]
pub deployment_msg: Option<String>,
}
#[derive(Deserialize, sqlx::FromRow)]
@@ -73,7 +76,6 @@ pub struct NewFlow {
pub schema: Option<Schema>,
pub draft_only: Option<bool>,
pub tag: Option<String>,
pub ws_error_handler_muted: Option<bool>,
pub dedicated_worker: Option<bool>,
pub timeout: Option<i32>,
pub deployment_message: Option<String>,
@@ -552,6 +554,7 @@ pub struct ListFlowQuery {
pub order_desc: Option<bool>,
pub starred_only: Option<bool>,
pub include_draft_only: Option<bool>,
pub with_deployment_msg: Option<bool>,
}
pub fn add_virtual_items_if_necessary(modules: &mut Vec<FlowModule>) {
+1
View File
@@ -318,6 +318,7 @@ pub enum JobPayload {
FlowDependencies {
path: String,
dedicated_worker: Option<bool>,
version: i64,
},
AppDependencies {
path: String,
+4
View File
@@ -206,6 +206,9 @@ pub struct ListableScript {
pub no_main_func: Option<bool>,
#[serde(skip_serializing_if = "is_false")]
pub use_codebase: bool,
#[sqlx(default)]
#[serde(skip_serializing_if = "Option::is_none")]
pub deployment_msg: Option<String>,
}
fn is_false(x: &bool) -> bool {
@@ -338,6 +341,7 @@ pub struct ListScriptQuery {
pub starred_only: Option<bool>,
pub include_without_main: Option<bool>,
pub include_draft_only: Option<bool>,
pub with_deployment_msg: Option<bool>,
}
pub fn to_i64(s: &str) -> crate::error::Result<i64> {
+1 -1
View File
@@ -18,7 +18,7 @@ pub type DB = Pool<Postgres>;
#[derive(Clone, Debug)]
pub enum DeployedObject {
Script { hash: ScriptHash, path: String, parent_path: Option<String> },
Flow { path: String, parent_path: Option<String> },
Flow { path: String, parent_path: Option<String>, version: i64 },
App { path: String, version: i64, parent_path: Option<String> },
Folder { path: String },
Resource { path: String, parent_path: Option<String> },
+13 -9
View File
@@ -2080,7 +2080,11 @@ pub async fn custom_concurrency_key(
async fn legacy_concurrency_key(db: &Pool<Postgres>, queued_job: &QueuedJob) -> Option<String> {
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),
@@ -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 {
@@ -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<R: rsmq_async::RsmqConnection + Send + S
"Cannot resolve flow dependencies for flow without path".to_string(),
)
})?;
let version = job
.script_hash
.clone()
.ok_or_else(|| {
Error::InternalErr("Flow Dependency requires script hash (flow version)".to_owned())
})?
.0;
let raw_flow = job.raw_flow.clone().map(|v| Ok(v)).unwrap_or_else(|| {
Err(Error::InternalErr(
"Flow Dependency requires raw flow".to_owned(),
@@ -547,6 +583,13 @@ pub async fn handle_flow_dependency_job<R: rsmq_async::RsmqConnection + Send + S
)
.execute(db)
.await?;
sqlx::query!(
"UPDATE flow_version SET value = $1 WHERE id = $2",
new_flow_value,
version
)
.execute(db)
.await?;
}
tx.commit().await?;
@@ -555,7 +598,7 @@ pub async fn handle_flow_dependency_job<R: rsmq_async::RsmqConnection + Send + S
&job.created_by,
&db,
&job.workspace_id,
DeployedObject::Flow { path: job_path, parent_path },
DeployedObject::Flow { path: job_path, parent_path, version },
deployment_message,
rsmq.clone(),
false,
@@ -1016,14 +1059,14 @@ pub async fn handle_app_dependency_job<R: rsmq_async::RsmqConnection + Send + Sy
) -> 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<R: rsmq_async::RsmqConnection + Send + Sy
)
.await?;
// Re-check cancelation to ensure we don't accidentially override a flow.
// Re-check cancelation to ensure we don't accidentially override an app.
if sqlx::query_scalar!("SELECT canceled FROM queue WHERE id = $1", job.id)
.fetch_optional(db)
.await
@@ -10,6 +10,8 @@
export let username: string
export let isConflict = false
let loading = false
let usernameInfo:
| {
username: string
@@ -42,26 +44,31 @@
const dispatch = createEventDispatcher()
async function renameUser() {
const automateUsernameCreation =
(await SettingService.getGlobal({ key: 'automate_username_creation' })) ?? false
loading = true
try {
const automateUsernameCreation =
(await SettingService.getGlobal({ key: 'automate_username_creation' })) ?? false
if (!automateUsernameCreation) {
sendUserToast(
'Modifying the username is only possible when the creation of usernames is automated and defined at instance level..'
)
return
}
await UserService.globalUserRename({
email,
requestBody: {
new_username: username
if (!automateUsernameCreation) {
sendUserToast(
'Modifying the username is only possible when the creation of usernames is automated and defined at instance level..'
)
return
}
})
sendUserToast(`Rename user ${email} to ${username}`)
await UserService.globalUserRename({
email,
requestBody: {
new_username: username
}
})
dispatch('renamed')
sendUserToast(`Renamed user ${email} to ${username}`)
dispatch('renamed')
} finally {
loading = false
}
}
</script>
@@ -139,6 +146,7 @@
})
}}
disabled={email === undefined || !username}
{loading}
>
Confirm username change
</Button>
@@ -0,0 +1,81 @@
<script lang="ts">
import { createPopperActions, type PopperOptions } from 'svelte-popperjs'
import type { PopoverPlacement } from './Popover.model'
import Portal from 'svelte-portal'
export let placement: PopoverPlacement = 'bottom-end'
export let notClickable = false
export let disablePopup = false
export let disappearTimeout = 100
export let appearTimeout = 300
export let style: string | undefined = undefined
export let focusEl: HTMLElement | undefined = undefined
const [popperRef, popperContent] = createPopperActions({ placement })
const popperOptions: PopperOptions<{}> = {
placement,
strategy: 'fixed',
modifiers: [
{ name: 'offset', options: { offset: [8, 8] } },
{
name: 'arrow',
options: {
padding: 10
}
}
]
}
let showTooltip = false
let timeout: NodeJS.Timeout | undefined = undefined
let inTimeout: NodeJS.Timeout | undefined = undefined
function open() {
clearTimeout(timeout)
if (appearTimeout == 0) {
showTooltip = true
} else {
inTimeout = setTimeout(() => (showTooltip = true), appearTimeout)
}
}
function close() {
inTimeout && clearTimeout(inTimeout)
inTimeout = undefined
timeout = setTimeout(() => (showTooltip = false), disappearTimeout)
}
$: focusEl && focusEl?.focus()
</script>
{#if notClickable}
<!-- svelte-ignore a11y-no-static-element-interactions -->
<span {style} use:popperRef on:mouseenter={open} on:mouseleave={close} class={$$props.class}>
<slot />
</span>
{:else}
<button
{style}
use:popperRef
on:mouseenter={open}
on:mouseleave={close}
on:click
class={$$props.class}
>
<slot />
</button>
{/if}
{#if showTooltip && !disablePopup}
<Portal>
<!-- svelte-ignore a11y-no-static-element-interactions -->
<div
use:popperContent={popperOptions}
on:mouseenter={open}
on:mouseleave={close}
class="z-[5001] border rounded-lg shadow-lg p-4 bg-surface"
>
<slot name="overlay" />
</div>
</Portal>
{/if}
+120 -18
View File
@@ -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<void> {
async function saveFlow(deploymentMsg?: string): Promise<void> {
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<Icon>
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
</script>
<svelte:window on:keydown={onKeyDown} />
@@ -998,6 +1042,10 @@
{#key renderCount}
{#if !$userStore?.operator}
<FlowCopilotDrawer {getHubCompletions} {genFlow} bind:flowCopilotMode />
{#if $pathStore}
<FlowHistory bind:this={flowHistory} path={$pathStore} on:historyRestore />
{/if}
<FlowImportExportMenu bind:drawer={jsonViewerDrawer} />
<FlowCopilotInputsModal
on:confirmed={async () => {
applyCopilotFlowInputs()
@@ -1078,10 +1126,40 @@
/>
</div>
</div>
<div class="flex flex-row space-x-2">
<div class="flex flex-row gap-2 items-center">
{#if $enterpriseLicense && !newFlow}
<Awareness />
{/if}
<div>
<ButtonDropdown hasPadding={false}>
<svelte:fragment slot="buttonReplacement">
<Button nonCaptureEvent size="xs" color="light">
<div class="flex flex-row items-center">
<MoreVertical size={14} />
</div>
</Button>
</svelte:fragment>
<svelte:fragment slot="items">
{#each moreItems as item}
<MenuItem
on:click={item.action}
disabled={item.disabled}
class={item.disabled ? 'opacity-50' : ''}
>
<div
class={twMerge(
'text-primary flex flex-row items-center text-left px-4 py-2 gap-2 cursor-pointer hover:bg-surface-hover !text-xs font-semibold'
)}
>
<svelte:component this={item.icon} size={14} />
{item.displayName}
</div>
</MenuItem>
{/each}
</svelte:fragment>
</ButtonDropdown>
</div>
<FlowBuilderTutorials
on:reload={() => {
renderCount += 1
@@ -1119,8 +1197,6 @@
{abortController}
/>
<FlowImportExportMenu />
<FlowPreviewButtons />
<Button
loading={loadingDraft}
@@ -1134,15 +1210,41 @@
>
Draft
</Button>
<Button
loading={loadingSave}
size="xs"
startIcon={{ icon: Save }}
on:click={() => saveFlow()}
dropdownItems={!newFlow ? dropdownItems : undefined}
>
Deploy
</Button>
<CustomPopover appearTimeout={0} focusEl={msgInput}>
<Button
loading={loadingSave}
size="xs"
startIcon={{ icon: Save }}
on:click={() => saveFlow()}
dropdownItems={!newFlow ? dropdownItems : undefined}
>
Deploy
</Button>
<svelte:fragment slot="overlay">
<div class="flex flex-row gap-2 w-80">
<input
type="text"
placeholder="Deployment message"
bind:value={deploymentMsg}
on:keydown={(e) => {
if (e.key === 'Enter') {
saveFlow(deploymentMsg)
}
}}
bind:this={msgInput}
/>
<Button
size="xs"
on:click={() => saveFlow(deploymentMsg)}
endIcon={{ icon: CornerDownLeft }}
loading={loadingSave}
>
Deploy
</Button>
</div>
</svelte:fragment>
</CustomPopover>
</div>
</div>
@@ -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
@@ -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<void> {
async function editScript(stay: boolean, deploymentMsg?: string): Promise<void> {
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
</script>
<svelte:window on:keydown={onKeyDown} />
@@ -1101,15 +1107,41 @@
>
<span class="hidden lg:flex"> Draft </span>
</Button>
<Button
loading={loadingSave}
size="xs"
startIcon={{ icon: Save }}
on:click={() => editScript(false)}
dropdownItems={computeDropdownItems(initialPath)}
>
Deploy
</Button>
<CustomPopover appearTimeout={0} focusEl={msgInput}>
<Button
loading={loadingSave}
size="xs"
startIcon={{ icon: Save }}
on:click={() => editScript(false)}
dropdownItems={computeDropdownItems(initialPath)}
>
Deploy
</Button>
<svelte:fragment slot="overlay">
<div class="flex flex-row gap-2 min-w-72">
<input
type="text"
placeholder="Deployment message"
bind:value={deploymentMsg}
bind:this={msgInput}
on:keydown={(e) => {
if (e.key === 'Enter') {
editScript(false, deploymentMsg)
}
}}
/>
<Button
size="xs"
on:click={() => editScript(false, deploymentMsg)}
endIcon={{ icon: CornerDownLeft }}
loading={loadingSave}
>
Deploy
</Button>
</div>
</svelte:fragment>
</CustomPopover>
</div>
</div>
</div>
@@ -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
}
@@ -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)
}
</script>
{#if menuOpen}
<AppJsonEditor on:change bind:this={appExport} />
<AppDeploymentHistory bind:this={appDeploymentHistory} />
<AppDeploymentHistory bind:this={appDeploymentHistory} appPath={app.path} />
{/if}
<Row
@@ -189,7 +180,7 @@
{
displayName: 'Deployments',
icon: History,
action: () => loadDeployements(),
action: () => appDeploymentHistory?.open(),
hide: $userStore?.operator
},
{
@@ -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
</script>
{#if menuOpen}
<ScheduleEditor on:update={() => goto('/schedules')} bind:this={scheduleEditor} />
<FlowHistory bind:this={flowHistory} path={flow.path} />
{/if}
<Row
@@ -194,6 +198,14 @@
disabled: archived,
hide: $userStore?.operator
},
{
displayName: 'Deployments',
icon: HistoryIcon,
action: () => {
flowHistory.open()
},
hide: $userStore?.operator
},
{
displayName: 'Schedule',
icon: Calendar,
@@ -0,0 +1,214 @@
<script lang="ts">
import { Pane, Splitpanes } from 'svelte-splitpanes'
import PanelSection from '../apps/editor/settingsPanel/common/PanelSection.svelte'
import { classNames, displayDate, emptyString, sendUserToast } from '$lib/utils'
import { type Flow, FlowService, type FlowVersion } from '$lib/gen'
import { workspaceStore } from '$lib/stores'
import { Skeleton } from '$lib/components/common'
import FlowViewer from '../FlowViewer.svelte'
import Drawer from '../common/drawer/Drawer.svelte'
import DrawerContent from '../common/drawer/DrawerContent.svelte'
import Button from '../common/button/Button.svelte'
import { ArrowRight, Pencil, X } from 'lucide-svelte'
import { createEventDispatcher } from 'svelte'
export let path: string
let drawer: Drawer
let loading: boolean = false
let versions: FlowVersion[] = []
let selectedVersion: FlowVersion | undefined = undefined
let selected: Flow | undefined = undefined
let deploymentMsgUpdateMode = false
let deploymentMsgUpdate: string | undefined = undefined
export function open() {
loadVersions()
drawer.openDrawer()
}
async function loadFlow(version: number) {
selected = await FlowService.getFlowVersion({
workspace: $workspaceStore!,
version,
path
})
}
async function loadVersions() {
loading = true
versions = await FlowService.getFlowHistory({
workspace: $workspaceStore!,
path: path
})
loading = false
}
async function updateDeploymentMsg(version: number | undefined) {
if (
selectedVersion === undefined ||
version === undefined ||
emptyString(deploymentMsgUpdate)
) {
return
}
await FlowService.updateFlowHistory({
workspace: $workspaceStore!,
version,
path,
requestBody: {
deployment_msg: deploymentMsgUpdate!
}
})
selectedVersion.deployment_msg = deploymentMsgUpdate
deploymentMsgUpdateMode = false
loadVersions()
}
const dispatch = createEventDispatcher()
async function restoreVersion(flow: Flow | undefined) {
if (!flow) return
await FlowService.updateFlow({
workspace: $workspaceStore!,
requestBody: {
...flow,
path
},
path
})
dispatch('historyRestore')
drawer?.closeDrawer()
sendUserToast('Flow restored from previous deployment')
}
loadVersions()
$: selectedVersion !== undefined && loadFlow(selectedVersion.id)
</script>
<Drawer bind:this={drawer} size="1200px">
<DrawerContent
on:close={() => {
drawer?.closeDrawer()
}}
>
<Splitpanes class="!overflow-visible">
<Pane size={20}>
<PanelSection title="Past Deployments">
<div class="flex flex-col gap-2 w-full">
{#if !loading}
{#if versions.length > 0}
<div class="flex gap-2 flex-col">
{#each versions ?? [] as version}
<!-- svelte-ignore a11y-click-events-have-key-events -->
<div
class={classNames(
'border flex gap-1 truncate justify-between flex-row w-full items-center p-2 rounded-md cursor-pointer hover:bg-blue-50 hover:text-blue-400',
selectedVersion?.id == version.id ? 'bg-blue-100 text-blue-600' : ''
)}
role="button"
tabindex="0"
on:click={() => {
selectedVersion = version
}}
>
<span class="text-xs truncate">
{#if emptyString(version.deployment_msg)}Version {version.id}{:else}{version.deployment_msg}{/if}
</span>
</div>
{/each}
</div>
{:else}
<div class="text-sm text-tertiary">No items</div>
{/if}
{:else}
<Skeleton layout={[[40], [40], [40], [40], [40]]} />
{/if}
</div>
</PanelSection>
</Pane>
<Pane size={80}>
<div class="h-full w-full overflow-auto">
{#if selectedVersion}
{#if selected}
<div class="px-2 flex flex-col gap-2">
<span class="flex flex-row text-sm p-1 text-tertiary">
{#if deploymentMsgUpdateMode}
<div class="flex w-full">
<input
type="text"
bind:value={deploymentMsgUpdate}
class="!w-auto grow"
on:click|stopPropagation={() => {}}
on:keydown|stopPropagation
on:keypress|stopPropagation={({ key }) => {
if (key === 'Enter') updateDeploymentMsg(selectedVersion?.id)
}}
/>
<Button
size="xs"
color="blue"
buttonType="button"
btnClasses="!p-1 !w-[34px] !ml-1"
aria-label="Save Deployment Message"
on:click={() => {
updateDeploymentMsg(selectedVersion?.id)
}}
>
<ArrowRight size={14} />
</Button>
<Button
size="xs"
color="light"
buttonType="button"
btnClasses="!p-1 !w-[34px] !ml-1"
aria-label="Abort"
on:click={() => {
deploymentMsgUpdateMode = false
deploymentMsgUpdate = undefined
}}
>
<X size={14} />
</Button>
</div>
{:else}
{#if selectedVersion.deployment_msg}
{selectedVersion.deployment_msg}
{:else}
Deployed {displayDate(selected.edited_at)} by {selected.edited_by}
{/if}
<button
on:click={() => {
deploymentMsgUpdate = selectedVersion?.deployment_msg
deploymentMsgUpdateMode = true
}}
title="Update commit message"
class="flex items-center px-1 rounded-sm hover:text-primary text-secondary h-5"
aria-label="Update commit message"
>
<Pencil size={14} />
</button>
{/if}
</span>
<div class="flex p-1 gap-2">
<Button size="xs" on:click={() => restoreVersion(selected)}
>Redeploy with that version
</Button>
</div>
<FlowViewer flow={selected} />
</div>
{:else}
<Skeleton layout={[[40]]} />
{/if}
{:else}
<div class="text-sm p-2 text-tertiary"
>Select a deployment version to see its details</div
>
{/if}
</div>
</Pane>
</Splitpanes>
</DrawerContent>
</Drawer>
@@ -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>('FlowEditorContext')
let jsonViewerDrawer: Drawer
export let drawer: Drawer | undefined
</script>
<Button
btnClasses="mr-2"
size="xs"
variant="border"
color="light"
on:click={() => jsonViewerDrawer.toggleDrawer()}
startIcon={{ icon: FileJson }}
>
Export
</Button>
<Drawer bind:this={jsonViewerDrawer} size="800px">
<DrawerContent title="OpenFlow" on:close={() => jsonViewerDrawer.toggleDrawer()}>
<Drawer bind:this={drawer} size="800px">
<DrawerContent title="OpenFlow" on:close={() => drawer?.toggleDrawer()}>
{#if $flowStore}
<FlowViewer flow={cleanInputs($flowStore)} tab="raw" />
{/if}
@@ -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++
}
</script>
<DiffDrawer bind:this={diffDrawer} {restoreDeployed} {restoreDraft} />
@@ -188,11 +201,7 @@
{#if app}
<div class="h-screen">
<AppEditor
on:restore={(e) => {
sendUserToast('App restored from previous deployment')
app = e.detail
redraw++
}}
on:restore={onRestore}
summary={app.summary}
app={app.value}
path={app.path}
@@ -205,6 +205,9 @@
const { path, selectedId } = e.detail
goto(`/flows/edit/${path}?selected=${selectedId}`)
}}
on:historyRestore={() => {
loadFlow()
}}
{flowStore}
{flowStateStore}
initialPath={$page.params.path}
@@ -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
</script>
<Skeleton
@@ -283,6 +292,9 @@
loadFlow()
}}
/>
{#if flow}
<FlowHistory bind:this={flowHistory} path={flow.path} on:historyRestore={loadFlow} />
{/if}
{#if flow}
<DetailPageLayout