backend: fix workflow as code (#5239)
* backend: improve `/get_job_update` after v2 * backend: insert missing `workflow_as_code_status` on completion also insert `flow_status` from so we can remove the query when `_duration` is above 500 * backend: fix workflow_as_code after v2 * backend: add `workflow_as_code` worker test
This commit is contained in:
67
backend/.sqlx/query-0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701.json
generated
Normal file
67
backend/.sqlx/query-0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701.json
generated
Normal file
@@ -0,0 +1,67 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT\n c.id IS NOT NULL AS completed,\n q.id IS NOT NULL AND q.running AS running,\n SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs,\n COALESCE(r.memory_peak, c.memory_peak) AS mem_peak,\n CASE\n -- flow step:\n WHEN flow_step_id IS NOT NULL THEN NULL\n -- completed:\n WHEN c.id IS NOT NULL THEN COALESCE(\n c.workflow_as_code_status || c.flow_status,\n c.workflow_as_code_status,\n c.flow_status\n )\n -- not completed:\n ELSE COALESCE(\n f.workflow_as_code_status || f.flow_status,\n f.workflow_as_code_status,\n f.flow_status\n )\n END AS \"flow_status: sqlx::types::Json<Box<RawValue>>\",\n job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset,\n created_by AS \"created_by!\",\n CASE WHEN $4::BOOLEAN THEN (\n SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc'\n ) END AS progress\n FROM v2_job j\n LEFT JOIN v2_job_queue q USING (id)\n LEFT JOIN v2_job_runtime r USING (id)\n LEFT JOIN v2_job_status f USING (id)\n LEFT JOIN v2_job_completed c USING (id)\n LEFT JOIN job_logs ON job_logs.job_id = $3\n WHERE j.workspace_id = $2 AND j.id = $3",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "completed",
|
||||
"type_info": "Bool"
|
||||
},
|
||||
{
|
||||
"ordinal": 1,
|
||||
"name": "running",
|
||||
"type_info": "Bool"
|
||||
},
|
||||
{
|
||||
"ordinal": 2,
|
||||
"name": "logs",
|
||||
"type_info": "Text"
|
||||
},
|
||||
{
|
||||
"ordinal": 3,
|
||||
"name": "mem_peak",
|
||||
"type_info": "Int4"
|
||||
},
|
||||
{
|
||||
"ordinal": 4,
|
||||
"name": "flow_status: sqlx::types::Json<Box<RawValue>>",
|
||||
"type_info": "Jsonb"
|
||||
},
|
||||
{
|
||||
"ordinal": 5,
|
||||
"name": "log_offset",
|
||||
"type_info": "Int4"
|
||||
},
|
||||
{
|
||||
"ordinal": 6,
|
||||
"name": "created_by!",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 7,
|
||||
"name": "progress",
|
||||
"type_info": "Int4"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Int4",
|
||||
"Text",
|
||||
"Uuid",
|
||||
"Bool"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
false,
|
||||
null
|
||||
]
|
||||
},
|
||||
"hash": "0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701"
|
||||
}
|
||||
31
backend/.sqlx/query-3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008.json
generated
Normal file
31
backend/.sqlx/query-3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008.json
generated
Normal file
@@ -0,0 +1,31 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "INSERT INTO v2_job_completed AS cj\n ( workspace_id\n , id\n , started_at\n , duration_ms\n , result\n , result_columns\n , canceled_by\n , canceled_reason\n , flow_status\n , workflow_as_code_status\n , memory_peak\n , status\n )\n SELECT q.workspace_id, q.id, started_at, COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(started_at, now()))))*1000), $3, $10, $5, $6,\n flow_status, workflow_as_code_status,\n $8, CASE WHEN $4::BOOL THEN 'canceled'::job_status\n WHEN $7::BOOL THEN 'skipped'::job_status\n WHEN $2::BOOL THEN 'success'::job_status\n ELSE 'failure'::job_status END AS status\n FROM v2_job_queue q LEFT JOIN v2_job_status USING (id) WHERE q.id = $1\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3 RETURNING duration_ms AS \"duration_ms!\"",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "duration_ms!",
|
||||
"type_info": "Int8"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Uuid",
|
||||
"Bool",
|
||||
"Jsonb",
|
||||
"Bool",
|
||||
"Varchar",
|
||||
"Text",
|
||||
"Bool",
|
||||
"Int4",
|
||||
"Int8",
|
||||
"TextArray"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008"
|
||||
}
|
||||
16
backend/.sqlx/query-3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082.json
generated
Normal file
16
backend/.sqlx/query-3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082.json
generated
Normal file
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "INSERT INTO v2_job_status (id, workflow_as_code_status)\n VALUES ($1, JSONB_SET('{}'::JSONB, array[$2], $3))\n ON CONFLICT (id) DO UPDATE SET workflow_as_code_status =\n COALESCE(EXCLUDED.workflow_as_code_status, '{}'::JSONB) || $3",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Uuid",
|
||||
"Text",
|
||||
"Jsonb"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082"
|
||||
}
|
||||
15
backend/.sqlx/query-45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3.json
generated
Normal file
15
backend/.sqlx/query-45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3.json
generated
Normal file
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE v2_job_status SET\n workflow_as_code_status = jsonb_set(\n jsonb_set(\n COALESCE(workflow_as_code_status, '{}'::jsonb),\n array[$1],\n COALESCE(workflow_as_code_status->$1, '{}'::jsonb)\n ),\n array[$1, 'started_at'],\n to_jsonb(now()::text)\n )\n WHERE id = $2",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3"
|
||||
}
|
||||
16
backend/.sqlx/query-56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5.json
generated
Normal file
16
backend/.sqlx/query-56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5.json
generated
Normal file
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE v2_job_status SET\n workflow_as_code_status = jsonb_set(\n jsonb_set(\n COALESCE(workflow_as_code_status, '{}'::jsonb),\n array[$1],\n COALESCE(workflow_as_code_status->$1, '{}'::jsonb)\n ),\n array[$1, 'duration_ms'],\n to_jsonb($2::bigint)\n )\n WHERE id = $3",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Int8",
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5"
|
||||
}
|
||||
22
backend/.sqlx/query-867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb.json
generated
Normal file
22
backend/.sqlx/query-867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb.json
generated
Normal file
@@ -0,0 +1,22 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT workflow_as_code_status FROM v2_job_completed WHERE id = $1",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "workflow_as_code_status",
|
||||
"type_info": "Jsonb"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
true
|
||||
]
|
||||
},
|
||||
"hash": "867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb"
|
||||
}
|
||||
12
backend/.sqlx/query-976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f.json
generated
Normal file
12
backend/.sqlx/query-976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f.json
generated
Normal file
@@ -0,0 +1,12 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM _sqlx_migrations WHERE version=20250201145630 OR version=20250201145631",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": []
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f"
|
||||
}
|
||||
14
backend/.sqlx/query-d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1.json
generated
Normal file
14
backend/.sqlx/query-d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1.json
generated
Normal file
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "INSERT INTO v2_job_status (id, workflow_as_code_status) VALUES ($1, '{}'::JSONB)\n ON CONFLICT (id) DO NOTHING",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1"
|
||||
}
|
||||
@@ -21,7 +21,7 @@ SELECT
|
||||
CASE WHEN j.trigger_kind = 'schedule'::job_trigger_kind THEN j.trigger END
|
||||
AS schedule_path,
|
||||
j.permissioned_as,
|
||||
s.flow_status,
|
||||
COALESCE(s.flow_status, s.workflow_as_code_status) AS flow_status,
|
||||
j.raw_flow,
|
||||
j.flow_step_id IS NOT NULL AS is_flow_step,
|
||||
j.script_lang AS language,
|
||||
@@ -21,7 +21,7 @@ SELECT
|
||||
CASE WHEN j.trigger_kind = 'schedule'::job_trigger_kind THEN j.trigger END
|
||||
AS schedule_path,
|
||||
j.permissioned_as,
|
||||
c.flow_status,
|
||||
COALESCE(c.flow_status, c.workflow_as_code_status) AS flow_status,
|
||||
j.raw_flow,
|
||||
j.flow_step_id IS NOT NULL AS is_flow_step,
|
||||
j.script_lang AS language,
|
||||
3
backend/tests/fixtures/base.sql
vendored
3
backend/tests/fixtures/base.sql
vendored
@@ -152,7 +152,8 @@ BEGIN
|
||||
-- v2_job_status:
|
||||
IF EXISTS(SELECT 1 FROM v2_job_status WHERE id = OLD.id) THEN
|
||||
SELECT * INTO job_status FROM v2_job_status WHERE id = OLD.id;
|
||||
IF job_status.flow_status::TEXT IS DISTINCT FROM OLD.__flow_status::TEXT THEN
|
||||
IF COALESCE(job_status.flow_status, job_status.workflow_as_code_status)::TEXT IS DISTINCT FROM OLD.__flow_status::TEXT
|
||||
THEN
|
||||
RAISE EXCEPTION 'flow_status mismatch';
|
||||
END IF;
|
||||
IF job_status.flow_leaf_jobs::TEXT IS DISTINCT FROM OLD.__leaf_jobs::TEXT THEN
|
||||
|
||||
@@ -3828,6 +3828,71 @@ async fn test_job_labels(db: Pool<Postgres>) {
|
||||
test(&["z", "a", "x"]).await;
|
||||
}
|
||||
|
||||
#[cfg(feature = "python")]
|
||||
const WORKFLOW_AS_CODE: &str = r#"
|
||||
from wmill import task
|
||||
|
||||
import pandas as pd
|
||||
import numpy as np
|
||||
|
||||
@task()
|
||||
def heavy_compute(n: int):
|
||||
df = pd.DataFrame(np.random.randn(100, 4), columns=list('ABCD'))
|
||||
return df.sum().sum()
|
||||
|
||||
@task
|
||||
def send_result(res: int, email: str):
|
||||
print(f"Sending result {res} to {email}")
|
||||
return "OK"
|
||||
|
||||
def main(n: int):
|
||||
l = []
|
||||
for i in range(n):
|
||||
l.append(heavy_compute(i))
|
||||
print(l)
|
||||
return [send_result(sum(l), "example@example.com"), n]
|
||||
"#;
|
||||
|
||||
#[cfg(feature = "python")]
|
||||
#[sqlx::test(fixtures("base", "hello"))]
|
||||
async fn test_workflow_as_code(db: Pool<Postgres>) {
|
||||
initialize_tracing().await;
|
||||
let server = ApiServer::start(db.clone()).await;
|
||||
let port = server.addr.port();
|
||||
|
||||
// workflow as code require at least 2 workers:
|
||||
let db = &db;
|
||||
in_test_worker(
|
||||
&db,
|
||||
async move {
|
||||
let job = RunJob::from(JobPayload::Code(RawCode {
|
||||
language: ScriptLang::Python3,
|
||||
content: WORKFLOW_AS_CODE.into(),
|
||||
..RawCode::default()
|
||||
}))
|
||||
.arg("n", json!(3))
|
||||
.run_until_complete(&db, port)
|
||||
.await;
|
||||
|
||||
assert_eq!(job.json_result().unwrap(), json!(["OK", 3]));
|
||||
let workflow_as_code_status = sqlx::query_scalar!(
|
||||
"SELECT workflow_as_code_status FROM v2_job_completed WHERE id = $1",
|
||||
job.id
|
||||
)
|
||||
.fetch_one(db)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
workflow_as_code_status.get("name"),
|
||||
Some(&json!("send_result"))
|
||||
);
|
||||
},
|
||||
port,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
async fn test_for_versions<F: Future<Output = ()>>(
|
||||
version_flags: impl Iterator<Item = Arc<RwLock<bool>>>,
|
||||
test: impl Fn() -> F,
|
||||
|
||||
@@ -188,6 +188,16 @@ pub async fn migrate(db: &DB) -> Result<(), Error> {
|
||||
tracing::info!("Could not remove sqlx migration with version=20250201145632: {err:#}");
|
||||
}
|
||||
|
||||
// New version of `v2_as_queue` and `v2_as_completed_job` VIEWs.
|
||||
if let Err(err) = sqlx::query!(
|
||||
"DELETE FROM _sqlx_migrations WHERE version=20250201145630 OR version=20250201145631"
|
||||
)
|
||||
.execute(db)
|
||||
.await
|
||||
{
|
||||
tracing::info!("Could not remove sqlx migration with version=[20250201145630, 20250201145631] : {err:#}");
|
||||
}
|
||||
|
||||
match sqlx::migrate!("../migrations")
|
||||
.run_direct(&mut custom_migrator)
|
||||
.await
|
||||
|
||||
@@ -3425,13 +3425,25 @@ pub async fn run_workflow_as_code(
|
||||
|
||||
if !wkflow_query.skip_update.unwrap_or(false) {
|
||||
sqlx::query!(
|
||||
"UPDATE v2_job_status SET flow_status = jsonb_set(COALESCE(flow_status, '{}'::jsonb), array[$1], jsonb_set(jsonb_set('{}'::jsonb, '{scheduled_for}', to_jsonb(now()::text)), '{name}', to_jsonb($3::text))) WHERE id = $2",
|
||||
uuid.to_string(),
|
||||
"INSERT INTO v2_job_status (id, workflow_as_code_status)
|
||||
VALUES ($1, JSONB_SET('{}'::JSONB, array[$2], $3))
|
||||
ON CONFLICT (id) DO UPDATE SET workflow_as_code_status =
|
||||
COALESCE(EXCLUDED.workflow_as_code_status, '{}'::JSONB) || $3",
|
||||
job_id,
|
||||
entrypoint
|
||||
).execute(&mut *tx).await?;
|
||||
uuid.to_string(),
|
||||
serde_json::json!({ "scheduled_for": Utc::now(), "name": entrypoint }),
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
} else {
|
||||
tracing::info!("Skipping update of flow status for job {job_id} in workspace {w_id}");
|
||||
sqlx::query!(
|
||||
"INSERT INTO v2_job_status (id, workflow_as_code_status) VALUES ($1, '{}'::JSONB)
|
||||
ON CONFLICT (id) DO NOTHING",
|
||||
job_id,
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
}
|
||||
|
||||
if *CLOUD_HOSTED {
|
||||
@@ -5159,114 +5171,72 @@ async fn get_log_file(Path((_w_id, file_p)): Path<(String, String)>) -> error::R
|
||||
)));
|
||||
}
|
||||
|
||||
#[derive(Deserialize, sqlx::FromRow)]
|
||||
pub struct JobUpdateRow {
|
||||
pub running: bool,
|
||||
pub logs: Option<String>,
|
||||
pub mem_peak: Option<i32>,
|
||||
pub flow_status: Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
|
||||
pub log_offset: Option<i32>,
|
||||
pub created_by: String,
|
||||
}
|
||||
async fn get_job_update(
|
||||
OptAuthed(opt_authed): OptAuthed,
|
||||
Extension(db): Extension<DB>,
|
||||
Path((w_id, job_id)): Path<(String, Uuid)>,
|
||||
Query(JobUpdateQuery { running, log_offset, get_progress }): Query<JobUpdateQuery>,
|
||||
) -> error::JsonResult<JobUpdate> {
|
||||
Query(JobUpdateQuery { log_offset, get_progress, .. }): Query<JobUpdateQuery>,
|
||||
) -> JsonResult<JobUpdate> {
|
||||
let record = sqlx::query!(
|
||||
"SELECT
|
||||
running AS \"running!\",
|
||||
substr(concat(coalesce(v2_as_queue.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) AS logs,
|
||||
mem_peak,
|
||||
CASE WHEN is_flow_step is true then NULL else flow_status END AS \"flow_status: sqlx::types::Json<Box<RawValue>>\",
|
||||
job_logs.log_offset + char_length(job_logs.logs) + 1 AS log_offset,
|
||||
created_by AS \"created_by!\"
|
||||
FROM v2_as_queue
|
||||
LEFT JOIN job_logs ON job_logs.job_id = v2_as_queue.id
|
||||
WHERE v2_as_queue.workspace_id = $2 AND v2_as_queue.id = $3",
|
||||
c.id IS NOT NULL AS completed,
|
||||
q.id IS NOT NULL AND q.running AS running,
|
||||
SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs,
|
||||
COALESCE(r.memory_peak, c.memory_peak) AS mem_peak,
|
||||
CASE
|
||||
-- flow step:
|
||||
WHEN flow_step_id IS NOT NULL THEN NULL
|
||||
-- completed:
|
||||
WHEN c.id IS NOT NULL THEN COALESCE(
|
||||
c.workflow_as_code_status || c.flow_status,
|
||||
c.workflow_as_code_status,
|
||||
c.flow_status
|
||||
)
|
||||
-- not completed:
|
||||
ELSE COALESCE(
|
||||
f.workflow_as_code_status || f.flow_status,
|
||||
f.workflow_as_code_status,
|
||||
f.flow_status
|
||||
)
|
||||
END AS \"flow_status: sqlx::types::Json<Box<RawValue>>\",
|
||||
job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset,
|
||||
created_by AS \"created_by!\",
|
||||
CASE WHEN $4::BOOLEAN THEN (
|
||||
SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc'
|
||||
) END AS progress
|
||||
FROM v2_job j
|
||||
LEFT JOIN v2_job_queue q USING (id)
|
||||
LEFT JOIN v2_job_runtime r USING (id)
|
||||
LEFT JOIN v2_job_status f USING (id)
|
||||
LEFT JOIN v2_job_completed c USING (id)
|
||||
LEFT JOIN job_logs ON job_logs.job_id = $3
|
||||
WHERE j.workspace_id = $2 AND j.id = $3",
|
||||
log_offset,
|
||||
&w_id,
|
||||
job_id
|
||||
job_id,
|
||||
get_progress.unwrap_or(false)
|
||||
)
|
||||
.fetch_optional(&db)
|
||||
.await?;
|
||||
.await?
|
||||
.ok_or_else(|| Error::NotFound(format!("Job not found: {}", job_id)))?;
|
||||
|
||||
let progress: Option<i32> = if get_progress == Some(true) {
|
||||
sqlx::query_scalar!(
|
||||
"SELECT scalar_int FROM job_stats WHERE workspace_id = $1 AND job_id = $2 AND metric_id = $3",
|
||||
&w_id,
|
||||
job_id,
|
||||
"progress_perc"
|
||||
)
|
||||
.fetch_optional(&db)
|
||||
.await?.and_then(|inner| inner)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
if let Some(record) = record {
|
||||
if opt_authed.is_none() && record.created_by != "anonymous" {
|
||||
return Err(Error::BadRequest(
|
||||
"As a non logged in user, you can only see jobs ran by anonymous users".to_string(),
|
||||
));
|
||||
}
|
||||
log_job_view(&db, opt_authed.as_ref(), &w_id, &job_id).await?;
|
||||
Ok(Json(JobUpdate {
|
||||
running: if !running && record.running {
|
||||
Some(true)
|
||||
} else {
|
||||
None
|
||||
},
|
||||
log_offset: record.log_offset,
|
||||
completed: None,
|
||||
new_logs: record.logs,
|
||||
mem_peak: record.mem_peak,
|
||||
progress,
|
||||
flow_status: record
|
||||
.flow_status
|
||||
.map(|x: sqlx::types::Json<Box<RawValue>>| x.0),
|
||||
}))
|
||||
} else {
|
||||
let record = sqlx::query!(
|
||||
"SELECT
|
||||
substr(concat(coalesce(v2_as_completed_job.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) AS logs,
|
||||
mem_peak,
|
||||
CASE WHEN is_flow_step is true then NULL else flow_status END AS \"flow_status: sqlx::types::Json<Box<RawValue>>\",
|
||||
job_logs.log_offset + char_length(job_logs.logs) + 1 AS log_offset,
|
||||
created_by AS \"created_by!\"
|
||||
FROM v2_as_completed_job
|
||||
LEFT JOIN job_logs ON job_logs.job_id = v2_as_completed_job.id
|
||||
WHERE v2_as_completed_job.workspace_id = $2 AND v2_as_completed_job.id = $3",
|
||||
log_offset,
|
||||
&w_id,
|
||||
job_id
|
||||
)
|
||||
.fetch_optional(&db)
|
||||
.await?;
|
||||
if let Some(record) = record {
|
||||
if opt_authed.is_none() && record.created_by != "anonymous" {
|
||||
return Err(Error::BadRequest(
|
||||
"As a non logged in user, you can only see jobs ran by anonymous users"
|
||||
.to_string(),
|
||||
));
|
||||
}
|
||||
log_job_view(&db, opt_authed.as_ref(), &w_id, &job_id).await?;
|
||||
Ok(Json(JobUpdate {
|
||||
running: Some(false),
|
||||
completed: Some(true),
|
||||
log_offset: record.log_offset,
|
||||
new_logs: record.logs,
|
||||
mem_peak: record.mem_peak,
|
||||
progress,
|
||||
flow_status: record
|
||||
.flow_status
|
||||
.map(|x: sqlx::types::Json<Box<RawValue>>| x.0),
|
||||
}))
|
||||
} else {
|
||||
Err(error::Error::NotFound(format!("Job not found: {}", job_id)))
|
||||
}
|
||||
if opt_authed.is_none() && record.created_by != "anonymous" {
|
||||
return Err(Error::BadRequest(
|
||||
"As a non logged in user, you can only see jobs ran by anonymous users".to_string(),
|
||||
));
|
||||
}
|
||||
log_job_view(&db, opt_authed.as_ref(), &w_id, &job_id).await?;
|
||||
Ok(Json(JobUpdate {
|
||||
running: record.running,
|
||||
completed: record.completed,
|
||||
log_offset: record.log_offset,
|
||||
new_logs: record.logs,
|
||||
mem_peak: record.mem_peak,
|
||||
progress: record.progress,
|
||||
flow_status: record
|
||||
.flow_status
|
||||
.map(|x: sqlx::types::Json<Box<RawValue>>| x.0),
|
||||
}))
|
||||
}
|
||||
|
||||
pub fn filter_list_completed_query(
|
||||
|
||||
@@ -353,7 +353,7 @@ pub enum JobPayload {
|
||||
Noop,
|
||||
}
|
||||
|
||||
#[derive(Clone, Serialize, Deserialize, Debug)]
|
||||
#[derive(Clone, Serialize, Deserialize, Debug, Default)]
|
||||
pub struct RawCode {
|
||||
pub content: String,
|
||||
pub path: Option<String>,
|
||||
|
||||
@@ -24,11 +24,12 @@ use serde::{ser::SerializeSeq, Deserialize, Deserializer, Serialize};
|
||||
|
||||
use crate::utils::StripPath;
|
||||
|
||||
#[derive(Serialize, Deserialize, Debug, PartialEq, Copy, Clone, Hash, Eq, sqlx::Type)]
|
||||
#[derive(Serialize, Deserialize, Debug, PartialEq, Copy, Clone, Hash, Eq, sqlx::Type, Default)]
|
||||
#[sqlx(type_name = "SCRIPT_LANG", rename_all = "lowercase")]
|
||||
#[serde(rename_all(serialize = "lowercase", deserialize = "lowercase"))]
|
||||
pub enum ScriptLang {
|
||||
Nativets,
|
||||
#[default]
|
||||
Deno,
|
||||
Python3,
|
||||
Go,
|
||||
|
||||
@@ -591,28 +591,28 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
, canceled_by
|
||||
, canceled_reason
|
||||
, flow_status
|
||||
, workflow_as_code_status
|
||||
, memory_peak
|
||||
, status
|
||||
)
|
||||
VALUES ($1, $2, $3, COALESCE($12::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($3, now()))))*1000), $5, $13, $7, $8, $9,\
|
||||
$11, CASE WHEN $6::BOOL THEN 'canceled'::job_status
|
||||
WHEN $10::BOOL THEN 'skipped'::job_status
|
||||
WHEN $4::BOOL THEN 'success'::job_status
|
||||
ELSE 'failure'::job_status END)
|
||||
ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $5 RETURNING duration_ms AS \"duration_ms!\"",
|
||||
/* $1 */ queued_job.workspace_id,
|
||||
/* $2 */ queued_job.id,
|
||||
/* $3 */ queued_job.started_at,
|
||||
/* $4 */ success,
|
||||
/* $5 */ result as Json<&T>,
|
||||
/* $6 */ canceled_by.is_some(),
|
||||
/* $7 */ canceled_by.clone().map(|cb| cb.username).flatten(),
|
||||
/* $8 */ canceled_by.clone().map(|cb| cb.reason).flatten(),
|
||||
/* $9 */ &queued_job.flow_status as &Option<Json<Box<RawValue>>>,
|
||||
/* $10 */ skipped,
|
||||
/* $11 */ if mem_peak > 0 { Some(mem_peak) } else { None },
|
||||
/* $12 */ duration,
|
||||
/* $13 */ result_columns as Option<&Vec<String>>,
|
||||
SELECT q.workspace_id, q.id, started_at, COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(started_at, now()))))*1000), $3, $10, $5, $6,
|
||||
flow_status, workflow_as_code_status,
|
||||
$8, CASE WHEN $4::BOOL THEN 'canceled'::job_status
|
||||
WHEN $7::BOOL THEN 'skipped'::job_status
|
||||
WHEN $2::BOOL THEN 'success'::job_status
|
||||
ELSE 'failure'::job_status END AS status
|
||||
FROM v2_job_queue q LEFT JOIN v2_job_status USING (id) WHERE q.id = $1
|
||||
ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3 RETURNING duration_ms AS \"duration_ms!\"",
|
||||
/* $1 */ queued_job.id,
|
||||
/* $2 */ success,
|
||||
/* $3 */ result as Json<&T>,
|
||||
/* $4 */ canceled_by.is_some(),
|
||||
/* $5 */ canceled_by.clone().map(|cb| cb.username).flatten(),
|
||||
/* $6 */ canceled_by.clone().map(|cb| cb.reason).flatten(),
|
||||
/* $7 */ skipped,
|
||||
/* $8 */ if mem_peak > 0 { Some(mem_peak) } else { None },
|
||||
/* $9 */ duration,
|
||||
/* $10 */ result_columns as Option<&Vec<String>>,
|
||||
)
|
||||
.fetch_one(&mut *tx)
|
||||
.await
|
||||
@@ -633,28 +633,14 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
}
|
||||
|
||||
if !queued_job.is_flow_step {
|
||||
if _duration > 500
|
||||
&& (queued_job.job_kind == JobKind::Script
|
||||
|| queued_job.job_kind == JobKind::Preview)
|
||||
{
|
||||
if let Err(e) = sqlx::query!(
|
||||
"UPDATE v2_job_completed SET flow_status = f.flow_status FROM v2_job_status f WHERE v2_job_completed.id = $1 AND f.id = $1 AND v2_job_completed.workspace_id = $2",
|
||||
&queued_job.id,
|
||||
&queued_job.workspace_id
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await {
|
||||
tracing::error!("Could not update job duration: {}", e);
|
||||
}
|
||||
}
|
||||
if let Some(parent_job) = queued_job.parent_job {
|
||||
if let Err(e) = sqlx::query_scalar!(
|
||||
let _ = sqlx::query_scalar!(
|
||||
"UPDATE v2_job_status SET
|
||||
flow_status = jsonb_set(
|
||||
workflow_as_code_status = jsonb_set(
|
||||
jsonb_set(
|
||||
COALESCE(flow_status, '{}'::jsonb),
|
||||
COALESCE(workflow_as_code_status, '{}'::jsonb),
|
||||
array[$1],
|
||||
COALESCE(flow_status->$1, '{}'::jsonb)
|
||||
COALESCE(workflow_as_code_status->$1, '{}'::jsonb)
|
||||
),
|
||||
array[$1, 'duration_ms'],
|
||||
to_jsonb($2::bigint)
|
||||
@@ -665,9 +651,11 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
parent_job
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await {
|
||||
tracing::error!("Could not update parent job flow_status: {}", e);
|
||||
}
|
||||
.await
|
||||
.inspect_err(|e| tracing::error!(
|
||||
"Could not update parent job `duration_ms` in workflow as code status: {}",
|
||||
e,
|
||||
));
|
||||
}
|
||||
}
|
||||
// tracing::error!("Added completed job {:#?}", queued_job);
|
||||
|
||||
@@ -2006,13 +2006,13 @@ async fn handle_queued_job(
|
||||
.warn_after_seconds(5)
|
||||
.await?;
|
||||
} else if let Some(parent_job) = job.parent_job {
|
||||
if let Err(e) = sqlx::query_scalar!(
|
||||
let _ = sqlx::query_scalar!(
|
||||
"UPDATE v2_job_status SET
|
||||
flow_status = jsonb_set(
|
||||
workflow_as_code_status = jsonb_set(
|
||||
jsonb_set(
|
||||
COALESCE(flow_status, '{}'::jsonb),
|
||||
COALESCE(workflow_as_code_status, '{}'::jsonb),
|
||||
array[$1],
|
||||
COALESCE(flow_status->$1, '{}'::jsonb)
|
||||
COALESCE(workflow_as_code_status->$1, '{}'::jsonb)
|
||||
),
|
||||
array[$1, 'started_at'],
|
||||
to_jsonb(now()::text)
|
||||
@@ -2024,9 +2024,12 @@ async fn handle_queued_job(
|
||||
.execute(db)
|
||||
.warn_after_seconds(5)
|
||||
.await
|
||||
{
|
||||
tracing::error!("Could not update parent job started_at flow_status: {}", e);
|
||||
}
|
||||
.inspect_err(|e| {
|
||||
tracing::error!(
|
||||
"Could not update parent job `started_at` in workflow as code status: {}",
|
||||
e
|
||||
)
|
||||
});
|
||||
}
|
||||
|
||||
let started = Instant::now();
|
||||
|
||||
Reference in New Issue
Block a user