fix(backend): fix error handler progress update

This commit is contained in:
Ruben Fiszel
2022-10-30 13:47:28 +01:00
parent e96c5ca670
commit b02521d7ff
8 changed files with 355 additions and 134 deletions
+210 -54
View File
@@ -1,8 +1,9 @@
use futures::{stream, Stream};
use serde_json::json;
use sqlx::{postgres::PgListener, query_scalar, types::Uuid, Pool, Postgres, Transaction};
use sqlx::{postgres::PgListener, types::Uuid, Pool, Postgres, Transaction};
use windmill_api::jobs::{CompletedJob, Job};
use windmill_common::{
flow_status::{FlowStatus, FlowStatusModule},
flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform},
scripts::ScriptLang,
DEFAULT_SLEEP_QUEUE,
@@ -107,6 +108,29 @@ impl ApiServer {
}
}
async fn _print_job(id: Uuid, db: &Pool<Postgres>) -> Result<(), anyhow::Error> {
tracing::info!(
"{:#?}",
get_job_by_id(db.begin().await?, "test-workspace", id)
.await?
.0
);
Ok(())
}
fn get_module(cjob: &CompletedJob, id: &str) -> Option<FlowStatusModule> {
cjob.flow_status.clone().and_then(|fs| {
find_module_in_vec(
serde_json::from_value::<FlowStatus>(fs).unwrap().modules,
id,
)
})
}
fn find_module_in_vec(modules: Vec<FlowStatusModule>, id: &str) -> Option<FlowStatusModule> {
modules.into_iter().find(|s| s.id() == id)
}
mod suspend_resume {
use futures::{Stream, StreamExt};
@@ -135,16 +159,6 @@ mod suspend_resume {
}
}
async fn _print_job(id: Uuid, db: &Pool<Postgres>) -> Result<(), anyhow::Error> {
tracing::info!(
"{:#?}",
get_job_by_id(db.begin().await?, "test-workspace", id)
.await?
.0
);
Ok(())
}
fn flow() -> FlowValue {
serde_json::from_value(serde_json::json!({
"modules": [{
@@ -273,7 +287,7 @@ mod suspend_resume {
server.close().await.unwrap();
let result = completed_job_result(flow, &db).await;
let result = completed_job(flow, &db).await.result.unwrap();
assert_eq!(
json!({
@@ -310,7 +324,9 @@ mod suspend_resume {
.arg("op", json!("cancel"))
.arg("port", json!(port))
.run_until_complete(&db, port)
.await;
.await
.result
.unwrap();
server.close().await.unwrap();
@@ -374,7 +390,7 @@ mod suspend_resume {
server.close().await.unwrap();
let result = completed_job_result(flow, &db).await;
let result = completed_job(flow, &db).await.result.unwrap();
assert_eq!(
json!({"error": "Job canceled: approval request disapproved by unknown" }),
@@ -528,7 +544,9 @@ def main(last, port):
.arg("items", json!(["unused", "unused", "unused"]))
.arg("port", json!(server.addr.port()))
.run_until_complete(&db, server.addr.port())
.await;
.await
.result
.unwrap();
assert_eq!(server.close().await, attempts);
assert_eq!(json!([3, 5, 7, 9]), result);
@@ -555,7 +573,9 @@ def main(last, port):
.arg("items", json!(["unused", "unused", "unused"]))
.arg("port", json!(server.addr.port()))
.run_until_complete(&db, server.addr.port())
.await;
.await
.result
.unwrap();
assert_eq!(server.close().await, attempts);
assert!(result["error"]
@@ -587,7 +607,9 @@ def main(last, port):
.arg("items", json!(["unused", "unused", "unused"]))
.arg("port", json!(server.addr.port()))
.run_until_complete(&db, server.addr.port())
.await;
.await
.result
.unwrap();
assert_eq!(server.close().await, attempts);
assert!(result["error"]
@@ -604,8 +626,9 @@ def main(last, port):
let value = serde_json::from_value(json!({
"modules": [{
"input_transform": { "port": { "type": "javascript", "expr": "flow_input.port" } },
"id": "a",
"value": {
"input_transform": { "port": { "type": "javascript", "expr": "flow_input.port" } },
"type": "rawscript",
"language": "python3",
"content": r#"
@@ -617,9 +640,9 @@ def main(port):
"retry": { "constant": { "attempts": 1, "seconds": 0 } },
}],
"failure_module": {
"input_transform": { "error": { "type": "javascript", "expr": "previous_result", },
"port": { "type": "javascript", "expr": "flow_input.port" } },
"value": {
"input_transform": { "error": { "type": "javascript", "expr": "previous_result", },
"port": { "type": "javascript", "expr": "flow_input.port" } },
"type": "rawscript",
"language": "python3",
"content": r#"
@@ -643,10 +666,16 @@ def main(error, port):
.into_iter()
.unzip::<_, _, Vec<_>, Vec<_>>();
let server = Server::start(responses).await;
let result = RunJob::from(JobPayload::RawFlow { value, path: None })
let cjob = RunJob::from(JobPayload::RawFlow { value, path: None })
.arg("port", json!(server.addr.port()))
.run_until_complete(&db, server.addr.port())
.await;
let result = cjob.result.clone().unwrap();
let failed_module = get_module(&cjob, "a").unwrap();
match failed_module {
FlowStatusModule::Failure { .. } => {}
_ => panic!("expected failure module"),
}
assert_eq!(server.close().await, attempts);
assert_eq!(
@@ -694,14 +723,18 @@ async fn test_iteration(db: Pool<Postgres>) {
let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None })
.arg("items", json!([]))
.run_until_complete(&db, server.addr.port())
.await;
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([]));
/* Don't actually test that this does 257 jobs or that will take forever. */
let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None })
.arg("items", json!((0..257).collect::<Vec<_>>()))
.run_until_complete(&db, server.addr.port())
.await;
.await
.result
.unwrap();
assert!(matches!(result, serde_json::Value::Object(_)));
assert!(result["error"]
.as_str()
@@ -751,11 +784,11 @@ impl RunJob {
}
/// push the job, spawn a worker, wait until the job is in completed_job
async fn run_until_complete(self, db: &Pool<Postgres>, port: u16) -> serde_json::Value {
async fn run_until_complete(self, db: &Pool<Postgres>, port: u16) -> CompletedJob {
let uuid = self.push(db).await;
let listener = listen_for_completed_jobs(db).await;
in_test_worker(db, listener.find(&uuid), port).await;
completed_job_result(uuid, db).await
completed_job(uuid, db).await
}
}
@@ -763,7 +796,7 @@ async fn run_job_in_new_worker_until_complete(
db: &Pool<Postgres>,
job: JobPayload,
port: u16,
) -> serde_json::Value {
) -> CompletedJob {
RunJob::from(job).run_until_complete(db, port).await
}
@@ -880,8 +913,8 @@ async fn listen_for_uuid_on(
}))
}
async fn completed_job_result(uuid: Uuid, db: &Pool<Postgres>) -> serde_json::Value {
query_scalar("SELECT result FROM completed_job WHERE id = $1")
async fn completed_job(uuid: Uuid, db: &Pool<Postgres>) -> CompletedJob {
sqlx::query_as::<_, CompletedJob>("SELECT * FROM completed_job WHERE id = $1")
.bind(uuid)
.fetch_one(db)
.await
@@ -974,7 +1007,10 @@ async fn test_deno_flow(db: Pool<Postgres>) {
for i in 0..50 {
println!("deno flow iteration: {}", i);
let result = run_job_in_new_worker_until_complete(&db, job.clone(), port).await;
let result = run_job_in_new_worker_until_complete(&db, job.clone(), port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([2, 4, 6]), "iteration: {}", i);
}
}
@@ -1129,7 +1165,10 @@ async fn test_deno_flow_same_worker(db: Pool<Postgres>) {
let job = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, job.clone(), server.addr.port()).await;
let result = run_job_in_new_worker_until_complete(&db, job.clone(), server.addr.port())
.await
.result
.unwrap();
assert_eq!(
result,
serde_json::json!("false 1,true 1,false 1,true 2,false 1,true 3,false 1,true 3")
@@ -1181,7 +1220,10 @@ async fn test_flow_result_by_id(db: Pool<Postgres>) {
.unwrap();
let job = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, job.clone(), port).await;
let result = run_job_in_new_worker_until_complete(&db, job.clone(), port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([[42]]));
}
@@ -1195,8 +1237,9 @@ async fn test_stop_after_if(db: Pool<Postgres>) {
let flow: FlowValue = serde_json::from_value(serde_json::json!({
"modules": [
{
"input_transforms": { "n": { "type": "javascript", "expr": "flow_input.n" } },
"id": "a",
"value": {
"input_transforms": { "n": { "type": "javascript", "expr": "flow_input.n" } },
"type": "rawscript",
"language": "python3",
"content": "def main(n): return n",
@@ -1207,8 +1250,9 @@ async fn test_stop_after_if(db: Pool<Postgres>) {
},
},
{
"input_transforms": { "n": { "type": "javascript", "expr": "previous_result" } },
"id": "b",
"value": {
"input_transforms": { "n": { "type": "javascript", "expr": "previous_result" } },
"type": "rawscript",
"language": "python3",
"content": "def main(n): return f'last step saw {n}'",
@@ -1222,16 +1266,78 @@ async fn test_stop_after_if(db: Pool<Postgres>) {
let result = RunJob::from(job.clone())
.arg("n", json!(123))
.run_until_complete(&db, port)
.await;
.await
.result
.unwrap();
assert_eq!(json!("last step saw 123"), result);
let result = RunJob::from(job.clone())
let cjob = RunJob::from(job.clone())
.arg("n", json!(-123))
.run_until_complete(&db, port)
.await;
let result = cjob.result.unwrap();
assert_eq!(json!(-123), result);
}
#[sqlx::test(fixtures("base"))]
async fn test_stop_after_if_nested(db: Pool<Postgres>) {
initialize_tracing().await;
// let server = ApiServer::start(db.clone()).await;
// let port = server.addr.port();
let port = 123;
let flow: FlowValue = serde_json::from_value(serde_json::json!({
"modules": [
{
"id": "a",
"value": {
"branches": [{"modules": [{
"id": "b",
"value": {
"input_transforms": { "n": { "type": "javascript", "expr": "flow_input.n" } },
"type": "rawscript",
"language": "python3",
"content": "def main(n): return n",
},
"stop_after_if": {
"expr": "result < 0",
"skip_if_stopped": false,
}}]}],
"type": "branchall"
},
},
{
"id": "c",
"value": {
"input_transforms": { "n": { "type": "javascript", "expr": "previous_result" } },
"type": "rawscript",
"language": "python3",
"content": "def main(n): return f'last step saw {n}'",
},
},
],
}))
.unwrap();
let job = JobPayload::RawFlow { value: flow, path: None };
let result = RunJob::from(job.clone())
.arg("n", json!(123))
.run_until_complete(&db, port)
.await
.result
.unwrap();
assert_eq!(json!("last step saw [123]"), result);
let cjob = RunJob::from(job.clone())
.arg("n", json!(-123))
.run_until_complete(&db, port)
.await;
let result = cjob.result.unwrap();
assert_eq!(json!([-123]), result);
}
#[sqlx::test(fixtures("base"))]
async fn test_python_flow(db: Pool<Postgres>) {
initialize_tracing().await;
@@ -1282,7 +1388,9 @@ async fn test_python_flow(db: Pool<Postgres>) {
JobPayload::RawFlow { value: flow.clone(), path: None },
port,
)
.await;
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([2, 4, 6]), "iteration: {i}");
}
@@ -1315,7 +1423,9 @@ async fn test_python_flow_2(db: Pool<Postgres>) {
JobPayload::RawFlow { value: flow.clone(), path: None },
port,
)
.await;
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!("Hello"), "iteration: {i}");
}
@@ -1344,7 +1454,9 @@ func main(derp string) (string, error) {
}))
.arg("derp", json!("world"))
.run_until_complete(&db, port)
.await;
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!("hello world"));
}
@@ -1363,7 +1475,10 @@ def main():
let job = JobPayload::Code(RawCode { content, path: None, language: ScriptLang::Python3 });
let result = run_job_in_new_worker_until_complete(&db, job, port).await;
let result = run_job_in_new_worker_until_complete(&db, job, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!("hello world"));
}
@@ -1385,7 +1500,10 @@ def main():
let job = JobPayload::Code(RawCode { content, path: None, language: ScriptLang::Python3 });
let result = run_job_in_new_worker_until_complete(&db, job, port).await;
let result = run_job_in_new_worker_until_complete(&db, job, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!(3));
}
@@ -1406,7 +1524,10 @@ def main():
let job = JobPayload::Code(RawCode { content, path: None, language: ScriptLang::Python3 });
let result = run_job_in_new_worker_until_complete(&db, job, port).await;
let result = run_job_in_new_worker_until_complete(&db, job, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!("test-workspace"));
}
@@ -1458,7 +1579,10 @@ async fn test_empty_loop(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!(0));
}
@@ -1497,7 +1621,10 @@ async fn test_empty_loop_2(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([]));
}
@@ -1548,7 +1675,10 @@ async fn test_step_after_loop(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!(9));
}
@@ -1611,7 +1741,10 @@ async fn test_branchone_simple(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([1, 2]));
}
@@ -1644,7 +1777,10 @@ async fn test_branchall_simple(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([[1, 2], [1, 3]]));
}
@@ -1677,7 +1813,10 @@ async fn test_branchall_skip_failure(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(
result,
@@ -1707,7 +1846,10 @@ async fn test_branchall_skip_failure(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(
result,
@@ -1763,7 +1905,10 @@ async fn test_branchone_nested(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(result, serde_json::json!([1, 2, 3]));
}
@@ -1813,7 +1958,10 @@ async fn test_branchall_nested(db: Pool<Postgres>) {
.unwrap();
let flow = JobPayload::RawFlow { value: flow, path: None };
let result = run_job_in_new_worker_until_complete(&db, flow, port).await;
let result = run_job_in_new_worker_until_complete(&db, flow, port)
.await
.result
.unwrap();
assert_eq!(
result,
@@ -1873,7 +2021,9 @@ async fn test_failure_module(db: Pool<Postgres>) {
let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None })
.arg("n", json!(0))
.run_until_complete(&db, port)
.await;
.await
.result
.unwrap();
assert!(result["from failure module"]["error"]
.as_str()
.unwrap()
@@ -1882,7 +2032,9 @@ async fn test_failure_module(db: Pool<Postgres>) {
let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None })
.arg("n", json!(1))
.run_until_complete(&db, port)
.await;
.await
.result
.unwrap();
assert!(result["from failure module"]["error"]
.as_str()
.unwrap()
@@ -1891,7 +2043,9 @@ async fn test_failure_module(db: Pool<Postgres>) {
let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None })
.arg("n", json!(2))
.run_until_complete(&db, port)
.await;
.await
.result
.unwrap();
assert!(result["from failure module"]["error"]
.as_str()
.unwrap()
@@ -1900,6 +2054,8 @@ async fn test_failure_module(db: Pool<Postgres>) {
let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None })
.arg("n", json!(3))
.run_until_complete(&db, port)
.await;
.await
.result
.unwrap();
assert_eq!(json!({ "l": [0, 1, 2] }), result);
}
+27 -27
View File
@@ -21,12 +21,12 @@ use sqlx::{query_scalar, types::Uuid, Postgres, Transaction};
use windmill_audit::{audit_log, ActionKind};
use windmill_common::{
error::{self, to_anyhow, Error},
flow_status::{Approval, FlowStatus, FlowStatusModule},
flows::FlowValue,
oauth2::HmacSha256,
scripts::{ScriptHash, ScriptLang},
users::owner_to_token_owner,
utils::{not_found_if_none, now_from_db, paginate, require_admin, Pagination, StripPath},
worker_flow::{Approval, FlowStatus, FlowStatusModule},
};
use windmill_queue::{get_queued_job, push, JobKind, JobPayload, QueuedJob, RawCode};
@@ -208,32 +208,32 @@ pub async fn get_job_by_id<'c>(
#[derive(Debug, sqlx::FromRow, Serialize)]
pub struct CompletedJob {
workspace_id: String,
id: Uuid,
parent_job: Option<Uuid>,
created_by: String,
created_at: chrono::DateTime<chrono::Utc>,
started_at: chrono::DateTime<chrono::Utc>,
duration_ms: i32,
success: bool,
script_hash: Option<ScriptHash>,
script_path: Option<String>,
args: Option<serde_json::Value>,
result: Option<serde_json::Value>,
logs: Option<String>,
deleted: bool,
raw_code: Option<String>,
canceled: bool,
canceled_by: Option<String>,
canceled_reason: Option<String>,
job_kind: JobKind,
schedule_path: Option<String>,
permissioned_as: String,
flow_status: Option<serde_json::Value>,
raw_flow: Option<serde_json::Value>,
is_flow_step: bool,
language: Option<ScriptLang>,
is_skipped: bool,
pub workspace_id: String,
pub id: Uuid,
pub parent_job: Option<Uuid>,
pub created_by: String,
pub created_at: chrono::DateTime<chrono::Utc>,
pub started_at: chrono::DateTime<chrono::Utc>,
pub duration_ms: i32,
pub success: bool,
pub script_hash: Option<ScriptHash>,
pub script_path: Option<String>,
pub args: Option<serde_json::Value>,
pub result: Option<serde_json::Value>,
pub logs: Option<String>,
pub deleted: bool,
pub raw_code: Option<String>,
pub canceled: bool,
pub canceled_by: Option<String>,
pub canceled_reason: Option<String>,
pub job_kind: JobKind,
pub schedule_path: Option<String>,
pub permissioned_as: String,
pub flow_status: Option<serde_json::Value>,
pub raw_flow: Option<serde_json::Value>,
pub is_flow_step: bool,
pub language: Option<ScriptLang>,
pub is_skipped: bool,
}
#[derive(Deserialize, Clone, Copy)]
@@ -100,6 +100,8 @@ pub enum FlowStatusModule {
flow_jobs: Option<Vec<Uuid>>,
#[serde(skip_serializing_if = "Option::is_none")]
branch_chosen: Option<BranchChosen>,
#[serde(default)]
#[serde(skip_serializing_if = "Vec::is_empty")]
approvers: Vec<Approval>,
},
Failure {
@@ -145,7 +147,13 @@ impl FlowStatus {
.iter()
.map(|m| FlowStatusModule::WaitingForPriorSteps { id: m.id.clone() })
.collect(),
failure_module: FlowStatusModule::WaitingForPriorSteps { id: "failure".to_string() },
failure_module: FlowStatusModule::WaitingForPriorSteps {
id: f
.failure_module
.as_ref()
.map(|x| x.id.clone())
.unwrap_or_else(|| "failure".to_string()),
},
retry: RetryStatus { fail_count: 0, previous_result: None, failed_jobs: vec![] },
}
}
+1 -1
View File
@@ -12,6 +12,7 @@ use error::Error;
pub mod error;
pub mod external_ip;
pub mod flow_status;
pub mod flows;
pub mod more_serde;
pub mod oauth2;
@@ -19,7 +20,6 @@ pub mod scripts;
pub mod users;
pub mod utils;
pub mod variables;
pub mod worker_flow;
#[cfg(feature = "tracing_init")]
pub mod tracing_init;
+2 -2
View File
@@ -16,10 +16,10 @@ use uuid::Uuid;
use windmill_audit::{audit_log, ActionKind};
use windmill_common::{
error::{self, to_anyhow, Error},
flow_status::{init_flow_status, FlowStatus, MAX_RETRY_ATTEMPTS, MAX_RETRY_INTERVAL},
flows::FlowValue,
scripts::{get_full_hub_script_by_path, HubScript, ScriptHash, ScriptLang},
utils::StripPath,
worker_flow::{init_flow_status, FlowStatus, MAX_RETRY_ATTEMPTS, MAX_RETRY_INTERVAL},
};
lazy_static::lazy_static! {
@@ -200,7 +200,7 @@ pub async fn delete_job(
.map_err(|e| Error::InternalErr(format!("Error during deletion of job {job_id}: {e}")))?
.unwrap_or(0)
== 1;
tracing::debug!("Job {job_id} deletion was achieved with success: {job_removed}");
tracing::debug!("Job {job_id} deleted: {job_removed}");
Ok(())
}
+1 -1
View File
@@ -238,7 +238,7 @@ async function resource(path) {{
})
.join(""),
);
tracing::debug!("{}", code);
// tracing::debug!("{}", code);
let global = context.execute_script("<anon>", &code)?;
let global = context.resolve_value(global).await?;
+2 -6
View File
@@ -14,6 +14,7 @@ use uuid::Uuid;
use windmill_common::{
error::{self, to_anyhow, Error},
scripts::{ScriptHash, ScriptLang},
utils::rd_string,
variables,
};
use windmill_queue::{
@@ -57,12 +58,7 @@ pub async fn create_token_for_owner<'c>(
username: &str,
) -> error::Result<(Transaction<'c, Postgres>, String)> {
// TODO: Bad implementation. We should not have access to this DB here.
use rand::prelude::*;
let token: String = rand::thread_rng()
.sample_iter(&rand::distributions::Alphanumeric)
.take(30)
.map(char::from)
.collect();
let token: String = rd_string(30);
let is_super_admin = username.contains('@')
&& sqlx::query_scalar!(
"SELECT super_admin FROM password WHERE email = $1",
+103 -42
View File
@@ -21,11 +21,11 @@ use tracing::instrument;
use uuid::Uuid;
use windmill_common::{
error::{self, to_anyhow, Error},
flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform, Retry, Suspend},
worker_flow::{
flow_status::{
Approval, BranchAllStatus, BranchChosen, FlowStatus, FlowStatusModule, RetryStatus,
MAX_RETRY_ATTEMPTS, MAX_RETRY_INTERVAL,
},
flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform, Retry, Suspend},
};
type DB = sqlx::Pool<sqlx::Postgres>;
@@ -81,6 +81,8 @@ pub async fn update_flow_status_after_job_completion(
.and_then(|i| old_status.modules.get(i))
.unwrap_or(&old_status.failure_module);
tracing::debug!("UPDATE FLOW STATUS 2: {module_index:#?} {module_status:#?} {old_status:#?} ");
let skip_loop_failures = if matches!(
module_status,
FlowStatusModule::InProgress { iterator: Some(_), .. }
@@ -105,7 +107,7 @@ pub async fn update_flow_status_after_job_completion(
let (step_counter, new_status) = match module_status {
FlowStatusModule::InProgress {
iterator: Some(windmill_common::worker_flow::Iterator { index, itered, .. }),
iterator: Some(windmill_common::flow_status::Iterator { index, itered, .. }),
..
} if (*index + 1 < itered.len() && (success || skip_loop_failures)) => {
(old_status.step, module_status.clone())
@@ -154,7 +156,37 @@ pub async fn update_flow_status_after_job_completion(
.unwrap_or(true);
let (stop_early, skip_if_stop_early) = if let Some(se) = stop_early_override {
sqlx::query!(
"
UPDATE queue
SET flow_status = JSONB_SET(
JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2),
ARRAY['step'], $3)
WHERE id = $4
",
old_status.step.to_string(),
json!(new_status),
json!(step_counter),
flow
)
.execute(&mut tx)
.await?;
(true, se)
} else if old_status.step >= old_status.modules.len() as i32 {
tracing::debug!("SET NEW STATUS: {new_status:#?} ");
sqlx::query!(
"
UPDATE queue
SET flow_status = JSONB_SET(flow_status, ARRAY['failure_module'], $1)
WHERE id = $2
",
json!(new_status),
flow
)
.execute(&mut tx)
.await?;
(false, false)
} else {
let (stop_early_expr, skip_if_stop_early) = sqlx::query_as::<
_,
@@ -172,8 +204,8 @@ pub async fn update_flow_status_after_job_completion(
",
)
.bind(old_status.step)
.bind(serde_json::json!(new_status))
.bind(serde_json::json!(step_counter))
.bind(json!(new_status))
.bind(json!(step_counter))
.bind(flow)
.fetch_one(&mut tx)
.await
@@ -280,7 +312,6 @@ pub async fn update_flow_status_after_job_completion(
} else {
"Flow job completed".to_string()
};
tracing::debug!("{skip_if_stop_early:?}");
if flow_job.canceled {
add_completed_job_error(
db,
@@ -465,31 +496,54 @@ pub async fn update_flow_status_in_progress(
job_in_progress: Uuid,
) -> error::Result<()> {
let step = get_step_of_flow_status(db, flow).await?;
sqlx::query(&format!(
"UPDATE queue
SET flow_status = jsonb_set(jsonb_set(flow_status, '{{modules, {step}, job}}', $1), '{{modules, {step}, type}}', $2)
WHERE id = $3 AND workspace_id = $4",
))
.bind(json!(job_in_progress.to_string()))
.bind(json!("InProgress"))
.bind(flow)
.bind(w_id)
.execute(db)
.await?;
if let Step::Step(step) = step {
sqlx::query(&format!(
"UPDATE queue
SET flow_status = jsonb_set(jsonb_set(flow_status, '{{modules, {step}, job}}', $1), '{{modules, {step}, type}}', $2)
WHERE id = $3 AND workspace_id = $4",
))
.bind(json!(job_in_progress.to_string()))
.bind(json!("InProgress"))
.bind(flow)
.bind(w_id)
.execute(db)
.await?;
} else {
sqlx::query(&format!(
"UPDATE queue
SET flow_status = jsonb_set(jsonb_set(flow_status, '{{failure_module, job}}', $1), '{{failure_module, type}}', $2)
WHERE id = $3 AND workspace_id = $4",
))
.bind(json!(job_in_progress.to_string()))
.bind(json!("InProgress"))
.bind(flow)
.bind(w_id)
.execute(db)
.await?;
}
Ok(())
}
pub enum Step {
Step(i32),
FailureStep,
}
#[instrument(level = "trace", skip_all)]
pub async fn get_step_of_flow_status(db: &DB, id: Uuid) -> error::Result<i32> {
let r = sqlx::query_scalar!(
"SELECT (flow_status->'step')::integer FROM queue WHERE id = $1",
pub async fn get_step_of_flow_status(db: &DB, id: Uuid) -> error::Result<Step> {
let r = sqlx::query!(
"SELECT (flow_status->'step')::integer as step, jsonb_array_length(flow_status->'modules') as len FROM queue WHERE id = $1",
id
)
.fetch_one(db)
.await
.map_err(|e| Error::InternalErr(format!("fetching step flow status: {e}")))?
.ok_or_else(|| Error::InternalErr(format!("not found step")))?;
Ok(r)
.map_err(|e| Error::InternalErr(format!("fetching step flow status: {e}")))?;
if r.step < r.len {
Ok(Step::Step(r.step.ok_or_else(|| {
Error::InternalErr("step is null".to_string())
})?))
} else {
Ok(Step::FailureStep)
}
}
/// resumes should be in order of timestamp ascending, so that more recent are at the end
@@ -520,7 +574,7 @@ async fn transform_input(
let flow_input = flow_args.clone().unwrap_or_else(|| json!({}));
let previous_result = flatten_previous_result(last_result.clone());
let context = vec![
("params".to_string(), serde_json::json!(mapped)),
("params".to_string(), json!(mapped)),
("previous_result".to_string(), previous_result),
("flow_input".to_string(), flow_input),
(
@@ -596,7 +650,7 @@ pub async fn handle_flow(
&Uuid::nil(),
flow_job.workspace_id.as_str(),
true,
serde_json::json!({}),
json!({}),
None,
true,
same_worker_tx,
@@ -613,7 +667,7 @@ pub async fn handle_flow(
serde_json::from_value::<FlowStatus>(flow_job.flow_status.clone().unwrap_or_default())
.with_context(|| format!("parse flow status {}", flow_job.id))?;
tracing::debug!("handle_flow: {:#?}", flow_job);
tracing::debug!("handle_flow: {:#?}", flow_job.flow_status);
push_next_flow_job(
flow_job,
status,
@@ -695,7 +749,7 @@ async fn push_next_flow_job(
.unwrap_or_else(|| status.failure_module.clone());
tracing::debug!(
"push_next_flow_job: module: {:#?}, status: {:#?}",
"push_next_flow_job {i}: module: {:#?}, status: {:#?}",
module.value,
status_module
);
@@ -1048,11 +1102,11 @@ async fn push_next_flow_job(
FlowStatusModule::InProgress {
job: uuid,
iterator: Some(windmill_common::worker_flow::Iterator { index, itered }),
iterator: Some(windmill_common::flow_status::Iterator { index, itered }),
flow_jobs: Some(flow_jobs),
branch_chosen: None,
branchall: None,
id: module.id.clone(),
id: status_module.id(),
}
}
NextStatus::NextBranchStep(NextBranch { mut flow_jobs, status, .. }) => {
@@ -1064,7 +1118,7 @@ async fn push_next_flow_job(
flow_jobs: Some(flow_jobs),
branch_chosen: None,
branchall: Some(status),
id: module.id.clone(),
id: status_module.id(),
}
}
@@ -1074,23 +1128,30 @@ async fn push_next_flow_job(
flow_jobs: None,
branch_chosen: Some(branch),
branchall: None,
id: module.id.clone(),
id: status_module.id(),
},
NextStatus::NextStep => {
FlowStatusModule::WaitingForExecutor { id: module.id.clone(), job: uuid }
FlowStatusModule::WaitingForExecutor { id: status_module.id(), job: uuid }
}
};
sqlx::query(
tracing::debug!("STATUS STEP: {:?} {i} {:#?}", status.step, new_status);
let json_pointer = if i >= flow.modules.len() {
"'failure_module'"
} else {
"'modules', $1::TEXT"
};
sqlx::query(&format!(
"
UPDATE queue
SET flow_status = JSONB_SET(
JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2),
ARRAY['step'], $3)
WHERE id = $4
",
)
.bind(status.step)
UPDATE queue
SET flow_status = JSONB_SET(
JSONB_SET(flow_status, ARRAY[{json_pointer}], $2),
ARRAY['step'], $3)
WHERE id = $4
"
))
.bind(i as i32)
.bind(json!(new_status))
.bind(json!(i))
.bind(flow_job.id)
@@ -1308,7 +1369,7 @@ async fn compute_next_flow_transform<'c>(
}
FlowStatusModule::InProgress {
iterator: Some(windmill_common::worker_flow::Iterator { itered, index }),
iterator: Some(windmill_common::flow_status::Iterator { itered, index }),
flow_jobs: Some(flow_jobs),
..
} => {