fix(backend): bash flow lock & add flow lock tests (#933)

* Fix Bash flow lock

* Tests w/ fixes

* Add Sequence privileges
This commit is contained in:
Kai Jellinghaus
2022-11-23 19:17:19 +01:00
committed by GitHub
parent e8d4cf2ba7
commit 4ddb3ec276
7 changed files with 159 additions and 32 deletions
@@ -0,0 +1 @@
-- Add down migration script here
@@ -0,0 +1,8 @@
-- Add up migration script here
GRANT ALL
ON ALL SEQUENCES IN SCHEMA public
TO windmill_user;
GRANT ALL
ON ALL SEQUENCES IN SCHEMA public
TO windmill_admin;
+3 -1
View File
@@ -11,6 +11,8 @@ INSERT INTO usr(workspace_id, email, username, is_admin, role) VALUES
INSERT INTO workspace_key(workspace_id, kind, key) VALUES
('test-workspace', 'cloud', 'test-key');
insert INTO token(token, email, label, super_admin) VALUES ('SECRET_TOKEN', 'test@windmill.dev', 'test token', true);
GRANT ALL PRIVILEGES ON TABLE workspace_key TO windmill_admin;
GRANT ALL PRIVILEGES ON TABLE workspace_key TO windmill_user;
@@ -45,4 +47,4 @@ EXECUTE FUNCTION "notify_queue" ();
AFTER UPDATE ON "queue"
FOR EACH ROW
WHEN (NEW.flow_status IS DISTINCT FROM OLD.flow_status)
EXECUTE FUNCTION "notify_queue" ();
EXECUTE FUNCTION "notify_queue" ();
+139 -20
View File
@@ -2330,24 +2330,143 @@ async fn test_failure_module(db: Pool<Postgres>) {
assert_eq!(json!({ "l": [0, 1, 2] }), result);
}
// #[cfg(test)]
// mod client_test {
// use windmill_common::error::to_anyhow;
#[sqlx::test(fixtures("base"))]
async fn test_flow_lock_all(db: Pool<Postgres>) {
use futures::StreamExt;
initialize_tracing().await;
let server = ApiServer::start(db.clone()).await;
let port = server.addr.port();
// #[tokio::test]
// async fn test_rust_client() -> Result<(), Box<dyn std::error::Error>> {
// println!(
// "{:#?}",
// windmill_api_client::create_client(
// "http://windmill.wimill.xyz",
// "XXXXXXXXXXXXXXX".to_string(),
// )
// .get_variable("demo", "u/ruben/test", Some(true))
// .await
// .map_err(to_anyhow)
// .map(|v| { v.into_inner() })?
// .value
// );
// Ok(())
// }
// }
let flow: windmill_api_client::types::OpenFlow = serde_json::from_value(serde_json::json!({
"summary": "",
"description": "",
"value": {
"modules": [
{
"id": "a",
"value": {
"lock": null,
"path": null,
"type": "rawscript",
"content": "import wmill\n\ndef main():\n return \"Test\"\n",
"language": "python3",
"input_transforms": {}
},
"summary": null,
"stop_after_if": null,
"input_transforms": {}
},
{
"id": "b",
"value": {
"lock": null,
"path": null,
"type": "rawscript",
"content": "import * as wmill from \"https://deno.land/x/windmill@v1.50.0/mod.ts\"\n\nexport async function main() {\n return \"Hello\"\n}\n",
"language": "deno",
"input_transforms": {}
},
"summary": null,
"stop_after_if": null,
"input_transforms": {}
},
{
"id": "c",
"value": {
"lock": null,
"path": null,
"type": "rawscript",
"content": "package inner\n\nimport (\n\t\"fmt\"\n\t\"rsc.io/quote\"\n wmill \"github.com/windmill-labs/windmill-go-client\"\n)\n\n// the main must return (interface{}, error)\n\nfunc main() (interface{}, error) {\n\tfmt.Println(\"Hello, World\")\n // v, _ := wmill.GetVariable(\"g/all/pretty_secret\")\n return \"Test\"\n}\n",
"language": "go",
"input_transforms": {}
},
"summary": null,
"stop_after_if": null,
"input_transforms": {}
},
{
"id": "d",
"value": {
"lock": null,
"path": null,
"type": "rawscript",
"content": "\n# the last line of the stdout is the return value\necho \"Hello $msg\"\n",
"language": "bash",
"input_transforms": {}
},
"summary": null,
"stop_after_if": null,
"input_transforms": {}
}
],
"failure_module": null
},
"schema": {
"type": "object",
"$schema": "https://json-schema.org/draft/2020-12/schema",
"required": [],
"properties": {}
}
}))
.unwrap();
let client = windmill_api_client::create_client(
&format!("http://localhost:{port}"),
"SECRET_TOKEN".to_owned(),
);
client
.create_flow(
"test-workspace",
&windmill_api_client::types::OpenFlowWPath {
open_flow: flow,
path: "g/all/flow_lock_all".to_owned(),
},
)
.await
.unwrap();
let mut str = listen_for_completed_jobs(&db).await;
let listen_first_job = str.next();
in_test_worker(&db, listen_first_job, port).await;
client
.get_flow_by_path("test-workspace", "g/all/flow_lock_all")
.await
.unwrap()
.into_inner()
.subtype_0
.value
.modules
.into_iter()
.for_each(|m| {
assert!(matches!(
m.value,
windmill_api_client::types::FlowModuleValue::Rawscript {
language: windmill_api_client::types::RawScriptLanguage::Deno | windmill_api_client::types::RawScriptLanguage::Bash,
lock: Some(ref lock),
..
} if lock == "")
|| matches!(
m.value,
windmill_api_client::types::FlowModuleValue::Rawscript {
language: windmill_api_client::types::RawScriptLanguage::Go | windmill_api_client::types::RawScriptLanguage::Python3,
lock: Some(ref lock),
..
} if lock.len() > 0)
);
});
}
#[sqlx::test(fixtures("base"))]
async fn test_rust_client(db: Pool<Postgres>) {
initialize_tracing().await;
let server = ApiServer::start(db.clone()).await;
let port = server.addr.port();
windmill_api_client::create_client(
&format!("http://localhost:{port}"),
"SECRET_TOKEN".to_string(),
)
.list_workspaces()
.await
.unwrap();
}
+4 -6
View File
@@ -131,7 +131,7 @@ paths:
/auth/logout:
post:
security: []
summary: logout
summary: logout
operationId: logout
tags:
- user
@@ -2386,7 +2386,7 @@ paths:
path:
type: string
value: {}
summary:
summary:
type: string
policy:
$ref: "#/components/schemas/Policy"
@@ -2474,7 +2474,7 @@ paths:
path:
type: string
args: {}
raw_code:
raw_code:
type: object
properties:
content:
@@ -2499,7 +2499,6 @@ paths:
schema:
type: string
/w/{workspace}/jobs/run/f/{path}:
post:
summary: run flow by path
@@ -4706,7 +4705,6 @@ components:
on_behalf_of:
type: string
ListableApp:
type: object
properties:
@@ -4749,7 +4747,7 @@ components:
type: string
format: date-time
value: {}
policy:
policy:
$ref: "#/components/schemas/Policy"
execution_mode:
type: string
+3 -2
View File
@@ -6,6 +6,7 @@
* LICENSE-AGPL for a copy of the license.
*/
use hyper::StatusCode;
use reqwest::Client;
use sql_builder::prelude::*;
@@ -137,7 +138,7 @@ async fn create_flow(
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Json(nf): Json<NewFlow>,
) -> Result<String> {
) -> Result<(StatusCode, String)> {
// cron::Schedule::from_str(&ns.schedule).map_err(|e| error::Error::BadRequest(e.to_string()))?;
let mut tx = user_db.clone().begin(&authed).await?;
@@ -200,7 +201,7 @@ async fn create_flow(
.await?;
tx.commit().await?;
Ok(nf.path.to_string())
Ok((StatusCode::CREATED, nf.path.to_string()))
}
async fn check_schedule_conflict<'c>(
+1 -3
View File
@@ -2016,9 +2016,7 @@ async fn capture_dependency_job(
ScriptLang::Deno => {
generate_deno_lock(job_id, job_raw_code, logs, job_dir, db, timeout, &envs).await
}
_ => Err(error::Error::InternalErr(
"Language incompatible with dep job".to_string(),
)),
ScriptLang::Bash => Ok("".to_owned()),
}
}