mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-08-19 00:02:03 +00:00
feat(backend): separate properly logs from result
This commit is contained in:
@@ -45,6 +45,8 @@ pub enum Error {
|
||||
HexErr(#[from] hex::FromHexError),
|
||||
#[error("Migrating database: {0}")]
|
||||
DatabaseMigration(#[from] MigrateError),
|
||||
#[error("Non-zero exit status: {0}")]
|
||||
ExitStatus(i32),
|
||||
#[error(transparent)]
|
||||
Anyhow(#[from] anyhow::Error),
|
||||
}
|
||||
|
||||
+1
-1
@@ -1485,7 +1485,7 @@ pub async fn add_completed_job_error<E: ToString + std::fmt::Debug>(
|
||||
false,
|
||||
false,
|
||||
serde_json::Value::Object(output_map.clone()),
|
||||
format!("\n{}\n{}", logs, e.to_string()),
|
||||
logs,
|
||||
)
|
||||
.await?;
|
||||
Ok((a, output_map))
|
||||
|
||||
+252
-335
@@ -7,13 +7,7 @@
|
||||
*/
|
||||
|
||||
use itertools::Itertools;
|
||||
use std::{
|
||||
borrow::Borrow,
|
||||
collections::HashMap,
|
||||
io, panic,
|
||||
process::{ExitStatus, Stdio},
|
||||
time::Duration,
|
||||
};
|
||||
use std::{borrow::Borrow, collections::HashMap, io, panic, process::Stdio, time::Duration};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::{
|
||||
@@ -410,7 +404,6 @@ async fn handle_queued_job(
|
||||
}
|
||||
_ => {
|
||||
let mut logs = "".to_string();
|
||||
let mut last_line = "{}".to_string();
|
||||
|
||||
if job.is_flow_step {
|
||||
update_flow_status_in_progress(
|
||||
@@ -423,21 +416,33 @@ async fn handle_queued_job(
|
||||
.await?;
|
||||
}
|
||||
|
||||
let execution = handle_job(
|
||||
&job,
|
||||
&job_dir,
|
||||
db,
|
||||
timeout,
|
||||
worker_name,
|
||||
worker_dir.clone(),
|
||||
&mut logs,
|
||||
&mut last_line,
|
||||
worker_config,
|
||||
envs,
|
||||
)
|
||||
.await;
|
||||
tracing::info!(
|
||||
worker = %worker_name,
|
||||
job_id = %job.id,
|
||||
workspace_id = %job.workspace_id,
|
||||
"handling job {}",
|
||||
job.id
|
||||
);
|
||||
|
||||
match execution {
|
||||
logs.push_str(&format!("job {} on worker {}\n", &job.id, &worker_name));
|
||||
|
||||
let result = if matches!(job.job_kind, JobKind::Dependencies) {
|
||||
handle_dependency_job(&job, &mut logs, job_dir, db, timeout, &envs).await
|
||||
} else {
|
||||
handle_code_execution_job(
|
||||
&job,
|
||||
db,
|
||||
job_dir,
|
||||
worker_dir,
|
||||
&mut logs,
|
||||
timeout,
|
||||
worker_config,
|
||||
envs,
|
||||
)
|
||||
.await
|
||||
};
|
||||
|
||||
match result {
|
||||
Ok(r) => {
|
||||
add_completed_job(db, &job, true, false, r.clone(), logs).await?;
|
||||
if job.is_flow_step {
|
||||
@@ -456,8 +461,26 @@ async fn handle_queued_job(
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
let (_, output_map) =
|
||||
add_completed_job_error(db, &job, logs, e, Some(metrics.clone())).await?;
|
||||
let error_message = match e {
|
||||
Error::ExitStatus(_) => {
|
||||
format!(
|
||||
"Error during execution of the script:\nlast 10 logs lines:\n{}",
|
||||
logs.lines()
|
||||
.skip(logs.lines().count().max(10) - 10)
|
||||
.join("\n")
|
||||
)
|
||||
}
|
||||
err @ _ => format!("error before termination: {err:#?}"),
|
||||
};
|
||||
|
||||
let (_, output_map) = add_completed_job_error(
|
||||
db,
|
||||
&job,
|
||||
logs,
|
||||
error_message,
|
||||
Some(metrics.clone()),
|
||||
)
|
||||
.await?;
|
||||
if job.is_flow_step {
|
||||
update_flow_status_after_job_completion(
|
||||
db,
|
||||
@@ -537,99 +560,16 @@ async fn transform_json_value(
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn handle_job(
|
||||
job: &QueuedJob,
|
||||
job_dir: &str,
|
||||
db: &DB,
|
||||
timeout: i32,
|
||||
worker_name: &str,
|
||||
worker_dir: &str,
|
||||
logs: &mut String,
|
||||
last_line: &mut String,
|
||||
worker_config: &WorkerConfig,
|
||||
envs: &Envs,
|
||||
) -> Result<serde_json::Value, Error> {
|
||||
tracing::info!(
|
||||
worker = %worker_name,
|
||||
job_id = %job.id,
|
||||
workspace_id = %job.workspace_id,
|
||||
"handling job {}",
|
||||
job.id
|
||||
);
|
||||
|
||||
logs.push_str(&format!("job {} on worker {}\n", &job.id, &worker_name));
|
||||
|
||||
let mut status: Result<ExitStatus, Error> =
|
||||
Err(Error::InternalErr("job not started".to_string()));
|
||||
|
||||
if matches!(job.job_kind, JobKind::Dependencies) {
|
||||
handle_dependency_job(
|
||||
job,
|
||||
logs,
|
||||
job_dir,
|
||||
&mut status,
|
||||
db,
|
||||
last_line,
|
||||
timeout,
|
||||
&envs,
|
||||
)
|
||||
.await?;
|
||||
} else {
|
||||
handle_code_execution_job(
|
||||
job,
|
||||
db,
|
||||
job_dir,
|
||||
worker_dir,
|
||||
logs,
|
||||
&mut status,
|
||||
last_line,
|
||||
timeout,
|
||||
worker_config,
|
||||
envs,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
|
||||
if status.is_ok() && status.as_ref().unwrap().success() {
|
||||
let result = serde_json::from_str::<serde_json::Value>(last_line).map_err(|e| {
|
||||
Error::ExecutionErr(format!(
|
||||
"result {} is not parsable.\n err: {}",
|
||||
last_line,
|
||||
e.to_string()
|
||||
))
|
||||
})?;
|
||||
Ok(result)
|
||||
} else {
|
||||
let err = match status {
|
||||
Ok(_) => {
|
||||
let s = format!(
|
||||
"Error during execution of the script\nlast 10 logs lines:\n{}",
|
||||
logs.lines()
|
||||
.skip(logs.lines().count().max(10) - 10)
|
||||
.join("\n")
|
||||
);
|
||||
logs.push_str("\n\n--- ERROR ---\n");
|
||||
s
|
||||
}
|
||||
Err(err) => format!("error before termination: {err}"),
|
||||
};
|
||||
Err(Error::ExecutionErr(err))
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_code_execution_job(
|
||||
job: &QueuedJob,
|
||||
db: &sqlx::Pool<sqlx::Postgres>,
|
||||
job_dir: &str,
|
||||
worker_dir: &str,
|
||||
logs: &mut String,
|
||||
status: &mut Result<ExitStatus, Error>,
|
||||
last_line: &mut String,
|
||||
timeout: i32,
|
||||
worker_config: &WorkerConfig,
|
||||
envs: &Envs,
|
||||
) -> Result<(), Error> {
|
||||
) -> error::Result<serde_json::Value> {
|
||||
let (inner_content, requirements_o, language) = if matches!(job.job_kind, JobKind::Preview)
|
||||
|| matches!(job.job_kind, JobKind::Script_Hub)
|
||||
{
|
||||
@@ -661,7 +601,6 @@ async fn handle_code_execution_job(
|
||||
worker_name = %worker_name,
|
||||
job_id = %job.id,
|
||||
workspace_id = %job.workspace_id,
|
||||
is_ok = status.is_ok(),
|
||||
"started {} job {}",
|
||||
&lang_str,
|
||||
job.id
|
||||
@@ -684,7 +623,7 @@ mount {{
|
||||
} else {
|
||||
"".to_string()
|
||||
};
|
||||
match language {
|
||||
let result = match language {
|
||||
None => {
|
||||
return Err(Error::ExecutionErr(
|
||||
"Require language to be not null".to_string(),
|
||||
@@ -700,14 +639,12 @@ mount {{
|
||||
worker_name,
|
||||
job,
|
||||
logs,
|
||||
status,
|
||||
db,
|
||||
last_line,
|
||||
timeout,
|
||||
&inner_content,
|
||||
&shared_mount,
|
||||
)
|
||||
.await?
|
||||
.await
|
||||
}
|
||||
Some(ScriptLang::Deno) => {
|
||||
handle_deno_job(
|
||||
@@ -719,11 +656,9 @@ mount {{
|
||||
job_dir,
|
||||
&inner_content,
|
||||
timeout,
|
||||
status,
|
||||
last_line,
|
||||
&shared_mount,
|
||||
)
|
||||
.await?;
|
||||
.await
|
||||
}
|
||||
Some(ScriptLang::Go) => {
|
||||
handle_go_job(
|
||||
@@ -736,23 +671,21 @@ mount {{
|
||||
timeout,
|
||||
job_dir,
|
||||
requirements_o,
|
||||
status,
|
||||
last_line,
|
||||
&shared_mount,
|
||||
)
|
||||
.await?
|
||||
.await
|
||||
}
|
||||
}
|
||||
};
|
||||
tracing::info!(
|
||||
worker_name = %worker_name,
|
||||
job_id = %job.id,
|
||||
workspace_id = %job.workspace_id,
|
||||
is_ok = status.is_ok(),
|
||||
is_ok = result.is_ok(),
|
||||
"finished {} job {}",
|
||||
&lang_str,
|
||||
job.id
|
||||
);
|
||||
Ok(())
|
||||
result
|
||||
}
|
||||
|
||||
async fn handle_go_job(
|
||||
@@ -765,10 +698,8 @@ async fn handle_go_job(
|
||||
timeout: i32,
|
||||
job_dir: &str,
|
||||
requirements_o: Option<String>,
|
||||
status: &mut Result<ExitStatus, Error>,
|
||||
last_line: &mut String,
|
||||
shared_mount: &str,
|
||||
) -> Result<(), Error> {
|
||||
) -> Result<serde_json::Value, Error> {
|
||||
//go does not like executing modules at temp root
|
||||
let job_dir = &format!("{job_dir}/go");
|
||||
if let Some(requirements) = requirements_o {
|
||||
@@ -789,9 +720,7 @@ async fn handle_go_job(
|
||||
inner_content,
|
||||
logs,
|
||||
job_dir,
|
||||
status,
|
||||
db,
|
||||
last_line,
|
||||
timeout,
|
||||
go_path,
|
||||
true,
|
||||
@@ -812,16 +741,8 @@ async fn handle_go_job(
|
||||
&job.created_by,
|
||||
)
|
||||
.await?;
|
||||
let args = if let Some(args) = &job.args {
|
||||
Some(
|
||||
transform_json_value(&token, &job.workspace_id, &base_internal_url, args.clone())
|
||||
.await?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let ser_args = serde_json::to_string(&args).map_err(|e| Error::ExecutionErr(e.to_string()))?;
|
||||
write_file(job_dir, "args.json", &ser_args).await?;
|
||||
create_args_and_out_file(job, &token, base_internal_url, job_dir).await?;
|
||||
|
||||
let spread = sig
|
||||
.args
|
||||
.into_iter()
|
||||
@@ -869,9 +790,16 @@ func main() {{
|
||||
fmt.Println(err)
|
||||
os.Exit(1)
|
||||
}}
|
||||
fmt.Println()
|
||||
fmt.Println("result:")
|
||||
fmt.Println(string(res_json))
|
||||
f, err := os.OpenFile("result.json", os.O_APPEND|os.O_WRONLY, os.ModeAppend)
|
||||
if err != nil {{
|
||||
fmt.Println(err)
|
||||
os.Exit(1)
|
||||
}}
|
||||
_, err = f.WriteString(string(res_json))
|
||||
if err != nil {{
|
||||
fmt.Println(err)
|
||||
os.Exit(1)
|
||||
}}
|
||||
}}
|
||||
|
||||
"#,
|
||||
@@ -924,8 +852,8 @@ func main() {{
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?
|
||||
};
|
||||
*status = handle_child(&job.id, db, logs, last_line, timeout, child).await;
|
||||
Ok(())
|
||||
handle_child(&job.id, db, logs, timeout, child).await?;
|
||||
read_result(job_dir).await
|
||||
}
|
||||
|
||||
async fn handle_deno_job(
|
||||
@@ -937,10 +865,8 @@ async fn handle_deno_job(
|
||||
job_dir: &str,
|
||||
inner_content: &String,
|
||||
timeout: i32,
|
||||
status: &mut Result<ExitStatus, Error>,
|
||||
last_line: &mut String,
|
||||
shared_mount: &str,
|
||||
) -> Result<(), Error> {
|
||||
) -> error::Result<serde_json::Value> {
|
||||
logs.push_str("\n\n--- DENO CODE EXECUTION ---\n");
|
||||
set_logs(logs, job.id, db).await;
|
||||
let _ = write_file(job_dir, "inner.ts", inner_content).await?;
|
||||
@@ -954,16 +880,7 @@ async fn handle_deno_job(
|
||||
&job.created_by,
|
||||
)
|
||||
.await?;
|
||||
let args = if let Some(args) = &job.args {
|
||||
Some(
|
||||
transform_json_value(&token, &job.workspace_id, &base_internal_url, args.clone())
|
||||
.await?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let ser_args = serde_json::to_string(&args).map_err(|e| Error::ExecutionErr(e.to_string()))?;
|
||||
write_file(job_dir, "args.json", &ser_args).await?;
|
||||
create_args_and_out_file(job, &token, base_internal_url, job_dir).await?;
|
||||
let spread = sig.args.into_iter().map(|x| x.name).join(",");
|
||||
let wrapper_content: String = format!(
|
||||
r#"
|
||||
@@ -976,9 +893,7 @@ const args = await Deno.readTextFile("args.json")
|
||||
async function run() {{
|
||||
let res: any = await main(...args);
|
||||
const res_json = JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value);
|
||||
console.log();
|
||||
console.log("result:");
|
||||
console.log(res_json);
|
||||
await Deno.writeTextFile("result.json", res_json);
|
||||
Deno.exit(0);
|
||||
}}
|
||||
run();
|
||||
@@ -1042,8 +957,27 @@ run();
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?
|
||||
};
|
||||
*status = handle_child(&job.id, db, logs, last_line, timeout, child).await;
|
||||
handle_child(&job.id, db, logs, timeout, child).await?;
|
||||
read_result(job_dir).await
|
||||
}
|
||||
|
||||
async fn create_args_and_out_file(
|
||||
job: &QueuedJob,
|
||||
token: &String,
|
||||
base_internal_url: &String,
|
||||
job_dir: &str,
|
||||
) -> Result<(), Error> {
|
||||
let args = if let Some(args) = &job.args {
|
||||
Some(
|
||||
transform_json_value(token, &job.workspace_id, &base_internal_url, args.clone())
|
||||
.await?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let ser_args = serde_json::to_string(&args).map_err(|e| Error::ExecutionErr(e.to_string()))?;
|
||||
write_file(job_dir, "args.json", &ser_args).await?;
|
||||
write_file(job_dir, "result.json", "").await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1064,13 +998,11 @@ async fn handle_python_job(
|
||||
worker_name: &str,
|
||||
job: &QueuedJob,
|
||||
logs: &mut String,
|
||||
status: &mut Result<ExitStatus, Error>,
|
||||
db: &sqlx::Pool<sqlx::Postgres>,
|
||||
last_line: &mut String,
|
||||
timeout: i32,
|
||||
inner_content: &String,
|
||||
shared_mount: &str,
|
||||
) -> Result<(), Error> {
|
||||
) -> error::Result<serde_json::Value> {
|
||||
let requirements =
|
||||
requirements_o.ok_or_else(|| Error::InternalErr(format!("lockfile missing")))?;
|
||||
|
||||
@@ -1150,72 +1082,62 @@ async fn handle_python_job(
|
||||
};
|
||||
|
||||
logs.push_str("\n--- PIP DEPENDENCIES INSTALL ---\n");
|
||||
*status = handle_child(&job.id, db, logs, last_line, timeout, child).await;
|
||||
let child = handle_child(&job.id, db, logs, timeout, child).await;
|
||||
tracing::info!(
|
||||
worker_name = %worker_name,
|
||||
job_id = %job.id,
|
||||
workspace_id = %job.workspace_id,
|
||||
is_ok = status.is_ok(),
|
||||
is_ok = child.is_ok(),
|
||||
"finished setting up python dependencies {}",
|
||||
job.id
|
||||
);
|
||||
child?;
|
||||
}
|
||||
if requirements.len() == 0 || status.is_ok() {
|
||||
logs.push_str("\n\n--- PYTHON CODE EXECUTION ---\n");
|
||||
logs.push_str("\n\n--- PYTHON CODE EXECUTION ---\n");
|
||||
|
||||
set_logs(logs, job.id, db).await;
|
||||
set_logs(logs, job.id, db).await;
|
||||
|
||||
let _ = write_file(job_dir, "inner.py", inner_content).await?;
|
||||
let _ = write_file(job_dir, "inner.py", inner_content).await?;
|
||||
|
||||
let sig = crate::parser_py::parse_python_signature(inner_content)?;
|
||||
let transforms = sig
|
||||
.args
|
||||
.into_iter()
|
||||
.map(|x| match x.typ {
|
||||
Typ::Bytes => {
|
||||
format!(
|
||||
"if \"{}\" in kwargs and kwargs[\"{}\"] is not None:\n \
|
||||
let sig = crate::parser_py::parse_python_signature(inner_content)?;
|
||||
let transforms = sig
|
||||
.args
|
||||
.into_iter()
|
||||
.map(|x| match x.typ {
|
||||
Typ::Bytes => {
|
||||
format!(
|
||||
"if \"{}\" in kwargs and kwargs[\"{}\"] is not None:\n \
|
||||
kwargs[\"{}\"] = base64.b64decode(kwargs[\"{}\"])\n",
|
||||
x.name, x.name, x.name, x.name
|
||||
)
|
||||
}
|
||||
Typ::Datetime => {
|
||||
format!(
|
||||
"if \"{}\" in kwargs and kwargs[\"{}\"] is not None:\n \
|
||||
x.name, x.name, x.name, x.name
|
||||
)
|
||||
}
|
||||
Typ::Datetime => {
|
||||
format!(
|
||||
"if \"{}\" in kwargs and kwargs[\"{}\"] is not None:\n \
|
||||
kwargs[\"{}\"] = datetime.strptime(kwargs[\"{}\"], \
|
||||
'%Y-%m-%dT%H:%M')\n",
|
||||
x.name, x.name, x.name, x.name
|
||||
)
|
||||
}
|
||||
_ => "".to_string(),
|
||||
})
|
||||
.collect::<Vec<String>>()
|
||||
.join("");
|
||||
x.name, x.name, x.name, x.name
|
||||
)
|
||||
}
|
||||
_ => "".to_string(),
|
||||
})
|
||||
.collect::<Vec<String>>()
|
||||
.join("");
|
||||
|
||||
let token = create_token_for_owner(
|
||||
&db,
|
||||
&job.workspace_id,
|
||||
&job.permissioned_as,
|
||||
"ephemeral-script",
|
||||
timeout * 2,
|
||||
&job.created_by,
|
||||
)
|
||||
.await?;
|
||||
let token = create_token_for_owner(
|
||||
&db,
|
||||
&job.workspace_id,
|
||||
&job.permissioned_as,
|
||||
"ephemeral-script",
|
||||
timeout * 2,
|
||||
&job.created_by,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let args = if let Some(args) = &job.args {
|
||||
Some(
|
||||
transform_json_value(&token, &job.workspace_id, &base_internal_url, args.clone())
|
||||
.await?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let ser_args =
|
||||
serde_json::to_string(&args).map_err(|e| Error::ExecutionErr(e.to_string()))?;
|
||||
write_file(job_dir, "args.json", &ser_args).await?;
|
||||
create_args_and_out_file(job, &token, base_internal_url, job_dir).await?;
|
||||
|
||||
let wrapper_content: String = format!(
|
||||
r#"
|
||||
let wrapper_content: String = format!(
|
||||
r#"
|
||||
import json
|
||||
import base64
|
||||
from datetime import datetime
|
||||
@@ -1230,69 +1152,67 @@ for k, v in list(kwargs.items()):
|
||||
{transforms}
|
||||
res = inner_script.main(**kwargs)
|
||||
res_json = json.dumps(res, separators=(',', ':'), default=str).replace('\n', '')
|
||||
print()
|
||||
print("result:")
|
||||
print(res_json)
|
||||
with open("result.json", 'w') as f:
|
||||
f.write(res_json)
|
||||
"#,
|
||||
);
|
||||
write_file(job_dir, "main.py", &wrapper_content).await?;
|
||||
);
|
||||
write_file(job_dir, "main.py", &wrapper_content).await?;
|
||||
|
||||
let mut reserved_variables = get_reserved_variables(job, &token, &base_url, db).await?;
|
||||
if !disable_nsjail {
|
||||
let _ = write_file(
|
||||
job_dir,
|
||||
"run.config.proto",
|
||||
&NSJAIL_CONFIG_RUN_PYTHON3_CONTENT
|
||||
.replace("{JOB_DIR}", job_dir)
|
||||
.replace("{CLONE_NEWUSER}", &(!disable_nuser).to_string())
|
||||
.replace("{SHARED_MOUNT}", shared_mount),
|
||||
)
|
||||
.await?;
|
||||
} else {
|
||||
reserved_variables.insert("PYTHONPATH".to_string(), format!("{job_dir}/dependencies"));
|
||||
}
|
||||
|
||||
tracing::info!(
|
||||
worker_name = %worker_name,
|
||||
job_id = %job.id,
|
||||
workspace_id = %job.workspace_id,
|
||||
"started python code execution {}",
|
||||
job.id
|
||||
);
|
||||
let child = if !disable_nsjail {
|
||||
Command::new(nsjail_path)
|
||||
.current_dir(job_dir)
|
||||
.env_clear()
|
||||
.envs(reserved_variables)
|
||||
.env("PATH", path_env)
|
||||
.env("BASE_INTERNAL_URL", base_internal_url)
|
||||
.args(vec![
|
||||
"--config",
|
||||
"run.config.proto",
|
||||
"--",
|
||||
python_path,
|
||||
"-u",
|
||||
"/tmp/main.py",
|
||||
])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?
|
||||
} else {
|
||||
Command::new(python_path)
|
||||
.current_dir(job_dir)
|
||||
.env_clear()
|
||||
.envs(reserved_variables)
|
||||
.env("PATH", path_env)
|
||||
.env("BASE_INTERNAL_URL", base_internal_url)
|
||||
.args(vec!["-u", "main.py"])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?
|
||||
};
|
||||
|
||||
*status = handle_child(&job.id, db, logs, last_line, timeout, child).await;
|
||||
let mut reserved_variables = get_reserved_variables(job, &token, &base_url, db).await?;
|
||||
if !disable_nsjail {
|
||||
let _ = write_file(
|
||||
job_dir,
|
||||
"run.config.proto",
|
||||
&NSJAIL_CONFIG_RUN_PYTHON3_CONTENT
|
||||
.replace("{JOB_DIR}", job_dir)
|
||||
.replace("{CLONE_NEWUSER}", &(!disable_nuser).to_string())
|
||||
.replace("{SHARED_MOUNT}", shared_mount),
|
||||
)
|
||||
.await?;
|
||||
} else {
|
||||
reserved_variables.insert("PYTHONPATH".to_string(), format!("{job_dir}/dependencies"));
|
||||
}
|
||||
Ok(())
|
||||
|
||||
tracing::info!(
|
||||
worker_name = %worker_name,
|
||||
job_id = %job.id,
|
||||
workspace_id = %job.workspace_id,
|
||||
"started python code execution {}",
|
||||
job.id
|
||||
);
|
||||
let child = if !disable_nsjail {
|
||||
Command::new(nsjail_path)
|
||||
.current_dir(job_dir)
|
||||
.env_clear()
|
||||
.envs(reserved_variables)
|
||||
.env("PATH", path_env)
|
||||
.env("BASE_INTERNAL_URL", base_internal_url)
|
||||
.args(vec![
|
||||
"--config",
|
||||
"run.config.proto",
|
||||
"--",
|
||||
python_path,
|
||||
"-u",
|
||||
"/tmp/main.py",
|
||||
])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?
|
||||
} else {
|
||||
Command::new(python_path)
|
||||
.current_dir(job_dir)
|
||||
.env_clear()
|
||||
.envs(reserved_variables)
|
||||
.env("PATH", path_env)
|
||||
.env("BASE_INTERNAL_URL", base_internal_url)
|
||||
.args(vec!["-u", "main.py"])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?
|
||||
};
|
||||
|
||||
handle_child(&job.id, db, logs, timeout, child).await?;
|
||||
read_result(job_dir).await
|
||||
}
|
||||
|
||||
async fn create_dependencies_dir(job_dir: &str) {
|
||||
@@ -1303,16 +1223,22 @@ async fn create_dependencies_dir(job_dir: &str) {
|
||||
.expect("could not create dependencies dir");
|
||||
}
|
||||
|
||||
async fn read_result(job_dir: &str) -> error::Result<serde_json::Value> {
|
||||
let mut file = File::open(format!("{job_dir}/result.json")).await?;
|
||||
let mut content = "".to_string();
|
||||
file.read_to_string(&mut content).await?;
|
||||
serde_json::from_str(&content)
|
||||
.map_err(|e| Error::ExecutionErr(format!("Error parsing result: {e}")))
|
||||
}
|
||||
|
||||
async fn handle_dependency_job(
|
||||
job: &QueuedJob,
|
||||
logs: &mut String,
|
||||
job_dir: &str,
|
||||
status: &mut error::Result<ExitStatus>,
|
||||
db: &sqlx::Pool<sqlx::Postgres>,
|
||||
last_line: &mut String,
|
||||
timeout: i32,
|
||||
Envs { go_path, pip_extra_index_url, pip_index_url, pip_trusted_host, .. }: &Envs,
|
||||
) -> error::Result<()> {
|
||||
) -> error::Result<serde_json::Value> {
|
||||
let content = match job.language {
|
||||
Some(ScriptLang::Python3) => {
|
||||
create_dependencies_dir(job_dir).await;
|
||||
@@ -1340,22 +1266,20 @@ async fn handle_dependency_job(
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?;
|
||||
*status = handle_child(&job.id, db, logs, last_line, timeout, child).await;
|
||||
if status.is_ok() && status.as_ref().unwrap().success() {
|
||||
let path_lock = format!("{job_dir}/requirements.txt");
|
||||
let mut file = File::open(path_lock).await?;
|
||||
handle_child(&job.id, db, logs, timeout, child)
|
||||
.await
|
||||
.map_err(|e| Error::ExecutionErr(format!("Lock file generation failed: {e:?}")))?;
|
||||
let path_lock = format!("{job_dir}/requirements.txt");
|
||||
let mut file = File::open(path_lock).await?;
|
||||
|
||||
let mut req_content = "".to_string();
|
||||
file.read_to_string(&mut req_content).await?;
|
||||
Ok(req_content
|
||||
.lines()
|
||||
.filter(|x| !x.trim_start().starts_with('#'))
|
||||
.map(|x| x.to_string())
|
||||
.collect::<Vec<String>>()
|
||||
.join("\n"))
|
||||
} else {
|
||||
Err(format!("Lock file generation failed: {status:?}"))
|
||||
}
|
||||
let mut req_content = "".to_string();
|
||||
file.read_to_string(&mut req_content).await?;
|
||||
Ok(req_content
|
||||
.lines()
|
||||
.filter(|x| !x.trim_start().starts_with('#'))
|
||||
.map(|x| x.to_string())
|
||||
.collect::<Vec<String>>()
|
||||
.join("\n"))
|
||||
}
|
||||
Some(ScriptLang::Go) => {
|
||||
let requirements = job
|
||||
@@ -1367,9 +1291,7 @@ async fn handle_dependency_job(
|
||||
&requirements,
|
||||
logs,
|
||||
job_dir,
|
||||
status,
|
||||
db,
|
||||
last_line,
|
||||
timeout,
|
||||
go_path,
|
||||
false,
|
||||
@@ -1382,11 +1304,6 @@ async fn handle_dependency_job(
|
||||
|
||||
match content {
|
||||
Ok(content) => {
|
||||
let as_json = json!(content);
|
||||
|
||||
*last_line =
|
||||
format!(r#"{{ "success": "Successful lock file generation", "lock": {as_json} }}"#);
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE script SET lock = $1 WHERE hash = $2 AND workspace_id = $3",
|
||||
&content,
|
||||
@@ -1395,6 +1312,7 @@ async fn handle_dependency_job(
|
||||
)
|
||||
.execute(db)
|
||||
.await?;
|
||||
Ok(json!({ "success": "Successful lock file generation", "lock": content }))
|
||||
}
|
||||
Err(error) => {
|
||||
sqlx::query!(
|
||||
@@ -1405,9 +1323,9 @@ async fn handle_dependency_job(
|
||||
)
|
||||
.execute(db)
|
||||
.await?;
|
||||
Err(Error::ExecutionErr(format!("Error locking file: {error}")))?
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn install_go_dependencies(
|
||||
@@ -1415,9 +1333,7 @@ async fn install_go_dependencies(
|
||||
code: &str,
|
||||
logs: &mut String,
|
||||
job_dir: &str,
|
||||
status: &mut Result<ExitStatus, Error>,
|
||||
db: &sqlx::Pool<sqlx::Postgres>,
|
||||
last_line: &mut String,
|
||||
timeout: i32,
|
||||
go_path: &str,
|
||||
preview: bool,
|
||||
@@ -1430,40 +1346,34 @@ async fn install_go_dependencies(
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?;
|
||||
|
||||
*status = handle_child(job_id, db, logs, last_line, timeout, child).await;
|
||||
if status.is_ok() {
|
||||
let child = Command::new(go_path)
|
||||
.current_dir(job_dir)
|
||||
.env("GOMEMLIMIT", "2000MiB")
|
||||
.args(vec!["mod", "tidy"])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?;
|
||||
*status = handle_child(job_id, db, logs, last_line, timeout, child).await;
|
||||
}
|
||||
if status.is_ok() && status.as_ref().unwrap().success() {
|
||||
if preview {
|
||||
Ok(String::new())
|
||||
} else {
|
||||
let mut req_content = "".to_string();
|
||||
handle_child(job_id, db, logs, timeout, child).await?;
|
||||
|
||||
let mut file = File::open(format!("{job_dir}/go.mod")).await?;
|
||||
file.read_to_string(&mut req_content).await?;
|
||||
let child = Command::new(go_path)
|
||||
.current_dir(job_dir)
|
||||
.env("GOMEMLIMIT", "2000MiB")
|
||||
.args(vec!["mod", "tidy"])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?;
|
||||
handle_child(job_id, db, logs, timeout, child)
|
||||
.await
|
||||
.map_err(|e| Error::ExecutionErr(format!("Lock file generation failed: {e:?}")))?;
|
||||
|
||||
req_content.push_str(&format!("\n{GO_REQ_SPLITTER}\n"));
|
||||
|
||||
if let Ok(mut file) = File::open(format!("{job_dir}/go.sum")).await {
|
||||
file.read_to_string(&mut req_content).await?;
|
||||
}
|
||||
|
||||
Ok(req_content)
|
||||
}
|
||||
if preview {
|
||||
Ok(String::new())
|
||||
} else {
|
||||
tracing::info!("go mod error");
|
||||
let mut req_content = "".to_string();
|
||||
|
||||
Err(error::Error::ExecutionErr(format!(
|
||||
"Lock file generation failed. Status: {status:?}",
|
||||
)))
|
||||
let mut file = File::open(format!("{job_dir}/go.mod")).await?;
|
||||
file.read_to_string(&mut req_content).await?;
|
||||
|
||||
req_content.push_str(&format!("\n{GO_REQ_SPLITTER}\n"));
|
||||
|
||||
if let Ok(mut file) = File::open(format!("{job_dir}/go.sum")).await {
|
||||
file.read_to_string(&mut req_content).await?;
|
||||
}
|
||||
|
||||
Ok(req_content)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1527,10 +1437,9 @@ async fn handle_child(
|
||||
job_id: &Uuid,
|
||||
db: &DB,
|
||||
logs: &mut String,
|
||||
last_line: &mut String,
|
||||
timeout: i32,
|
||||
mut child: Child,
|
||||
) -> crate::error::Result<ExitStatus> {
|
||||
) -> error::Result<()> {
|
||||
let timeout = Duration::from_secs(u64::try_from(timeout).expect("invalid timeout"));
|
||||
let ping_interval = Duration::from_secs(5);
|
||||
let cancel_check_interval = Duration::from_millis(500);
|
||||
@@ -1636,8 +1545,6 @@ async fn handle_child(
|
||||
Ok(line) => {
|
||||
append_with_limit(&mut joined, &line, &mut log_remaining);
|
||||
|
||||
*last_line = line;
|
||||
|
||||
if log_remaining == 0 {
|
||||
tracing::info!(%job_id, "Too many logs lines for job {job_id}");
|
||||
let _ = set_too_many_logs.send(true);
|
||||
@@ -1723,7 +1630,17 @@ async fn handle_child(
|
||||
_ if *too_many_logs.borrow() => Err(Error::ExecutionErr(
|
||||
"logs or result reached limit".to_string(),
|
||||
)),
|
||||
Ok(Ok(status)) => Ok(status),
|
||||
Ok(Ok(status)) => {
|
||||
if status.success() {
|
||||
Ok(())
|
||||
} else if let Some(code) = status.code() {
|
||||
Err(error::Error::ExitStatus(code))
|
||||
} else {
|
||||
Err(error::Error::ExecutionErr(
|
||||
"process terminated by signal".to_string(),
|
||||
))
|
||||
}
|
||||
}
|
||||
Ok(Err(kill_reason)) => Err(Error::ExecutionErr(format!(
|
||||
"job process killed because {kill_reason:#?}"
|
||||
))),
|
||||
@@ -3075,7 +2992,7 @@ def main(error, port):
|
||||
json!({
|
||||
"recv": 42,
|
||||
"from failure module": {
|
||||
"error": "Error during execution of the script\nlast 10 logs lines:\n\n\n--- PYTHON CODE EXECUTION ---\n\nTraceback (most recent call last):\n File \"/tmp/main.py\", line 14, in <module>\n res = inner_script.main(**kwargs)\n File \"/tmp/inner.py\", line 5, in main\n return sock.recv(1)[0]\nIndexError: index out of range",
|
||||
"error": "Error during execution of the script:\nlast 10 logs lines:\n\n\n--- PYTHON CODE EXECUTION ---\n\nTraceback (most recent call last):\n File \"/tmp/main.py\", line 14, in <module>\n res = inner_script.main(**kwargs)\n File \"/tmp/inner.py\", line 5, in main\n return sock.recv(1)[0]\nIndexError: index out of range",
|
||||
}
|
||||
})
|
||||
);
|
||||
|
||||
@@ -65,7 +65,7 @@
|
||||
</script>
|
||||
|
||||
<div class="inline-highlight">
|
||||
{#if result}
|
||||
{#if result != undefined}
|
||||
{#if resultKind && resultKind != 'json'}
|
||||
<div class="mb-2 text-gray-500 text-sm bg-gray-50/20">
|
||||
as JSON <input type="checkbox" bind:checked={forceJson} /></div
|
||||
|
||||
@@ -102,10 +102,10 @@
|
||||
</top>
|
||||
<down slot="down">
|
||||
<pre class="overflow-x-auto break-all relative h-full p-2 text-sm"
|
||||
>{#if testJob && 'result' in testJob && testJob.result != undefined}<DisplayResult
|
||||
>{#if testJob != undefined && 'result' in testJob && testJob.result != undefined}<DisplayResult
|
||||
result={testJob.result}
|
||||
/>
|
||||
{:else if testIsLoading}Waiting for Result...
|
||||
{:else if testIsLoading}Waiting for result...
|
||||
{:else}Test to see the result here
|
||||
{/if}
|
||||
</pre>
|
||||
|
||||
@@ -64,7 +64,10 @@
|
||||
|
||||
function onClick(event: MouseEvent) {
|
||||
dispatch('click', event)
|
||||
if (href) goto(href)
|
||||
if (href) {
|
||||
event.preventDefault()
|
||||
goto(href)
|
||||
}
|
||||
}
|
||||
|
||||
$: isSmall = size === 'xs' || size === 'sm'
|
||||
|
||||
@@ -76,11 +76,11 @@
|
||||
<LogViewer content={previewJob?.logs} isLoading={previewIsLoading} />
|
||||
</top>
|
||||
<down slot="down">
|
||||
<pre
|
||||
class="overflow-x-auto break-all relative h-full p-2 text-sm">{#if previewJob && 'result' in previewJob && previewJob.result}<DisplayResult
|
||||
<pre class="overflow-x-auto break-all relative h-full p-2 text-sm"
|
||||
>{#if previewJob != undefined && 'result' in previewJob && previewJob.result != undefined}<DisplayResult
|
||||
result={previewJob.result}
|
||||
/>
|
||||
{:else if previewIsLoading}Waiting for Result...
|
||||
{:else if previewIsLoading}Waiting for result...
|
||||
{:else}Test to see the result here
|
||||
{/if}
|
||||
</pre>
|
||||
|
||||
@@ -80,7 +80,7 @@
|
||||
|
||||
// If we get results, focus on that tab. Else, focus on logs
|
||||
function initView(): void {
|
||||
if (job && 'result' in job && job.result) {
|
||||
if (job && 'result' in job && job.result != undefined) {
|
||||
viewTab = 'result'
|
||||
} else if (viewTab == 'result') {
|
||||
viewTab = 'logs'
|
||||
@@ -389,7 +389,9 @@
|
||||
<HighlightCode language={job.language} code={job.raw_code} />
|
||||
{:else if job}No code is available
|
||||
{:else}Loading...{/if}
|
||||
{:else if job && 'result' in job && job.result}<DisplayResult result={job.result} />
|
||||
{:else if job != undefined && 'result' in job && job.result != undefined}<DisplayResult
|
||||
result={job.result}
|
||||
/>
|
||||
{:else if job}No output is available yet
|
||||
{:else}Loading...
|
||||
{/if}
|
||||
|
||||
@@ -79,6 +79,13 @@ mount {
|
||||
is_bind: true
|
||||
}
|
||||
|
||||
mount {
|
||||
src: "{JOB_DIR}/result.json"
|
||||
dst: "/tmp/result.json"
|
||||
is_bind: true
|
||||
rw: true
|
||||
}
|
||||
|
||||
mount {
|
||||
src: "/etc/ssl"
|
||||
dst: "/etc/ssl"
|
||||
|
||||
@@ -90,6 +90,13 @@ mount {
|
||||
is_bind: true
|
||||
}
|
||||
|
||||
mount {
|
||||
src: "{JOB_DIR}/result.json"
|
||||
dst: "/tmp/go/result.json"
|
||||
rw: true
|
||||
is_bind: true
|
||||
}
|
||||
|
||||
|
||||
mount {
|
||||
src: "/etc/ssl"
|
||||
|
||||
@@ -84,6 +84,14 @@ mount {
|
||||
is_bind: true
|
||||
}
|
||||
|
||||
mount {
|
||||
src: "{JOB_DIR}/result.json"
|
||||
dst: "/tmp/result.json"
|
||||
rw: true
|
||||
is_bind: true
|
||||
}
|
||||
|
||||
|
||||
mount {
|
||||
src: "{JOB_DIR}/dependencies"
|
||||
dst: "/tmp/dependencies"
|
||||
|
||||
Reference in New Issue
Block a user