Files
windmill/backend/windmill-worker/src/nu_executor.rs
T
+1 82620098d2 feat: add native result streaming (#6242)
* feat: add stream output feature to SSE job updates

Adds stream_output field to JobUpdate struct that extracts log lines
starting with '[wm_stream]:' from job logs. Regular logs now exclude
stream lines, which are captured separately for specialized handling.

- Added stream_output: Option<String> field to JobUpdate struct
- Created extract_stream_output_from_logs() function to filter stream lines
- Modified get_job_update_data() to use stream extraction logic
- SSE clients now receive both new_logs and stream_output in job updates

Co-authored-by: Ruben Fiszel <rubenfiszel@users.noreply.github.com>

* feat: rename stream_output to stream and handle newlines as \n

- Renamed `stream_output` field to `stream` in JobUpdate struct
- Updated extract_stream_output_from_logs to extract_stream_from_logs
- Changed stream output to join with literal \n instead of actual newlines
- Stream lines are properly excluded from regular new_logs

🤖 Generated with [Claude Code](https://claude.ai/code)

Co-Authored-By: Ruben Fiszel <rubenfiszel@users.noreply.github.com>

* decision tree nits

* push ee ref

* push ee ref

* fix: fix id renaming in apps

* remove duplicate caching (#6285)

* feat: migrate audit log ids to bigints (blocking migration for EE)

* fix(mcp): add proper check for mcp routes (#6282)

* add proper check for mcp routes

* cleaner

* apply to flow

* fix add checks scopes

---------

Co-authored-by: dieriba <dieriba.pro@gmail.com>

* chore(main): release 1.514.0 (#6283)

* chore(main): release 1.514.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>

* fix: pin tokio to 1.46.1 and aws-sdks-ts

* pin rustls to 0.23.29 + pin aws-sdk

* chore(main): release 1.514.1 (#6288)

* chore(main): release 1.514.1

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>

* fix: improve docker logs collection in docker mode

* support $res: string in form inputs of arrays

* fix import nit

* fix: fix DynSelect

* nits

* fix: resource-type-ts-parser (#6289)

* fix: resource types as arg in typescript handle imported defined types

* Update nix flake (#6291)

* merge

* Small UI fixes (#6294)

* fix step history not refreshing with staticInputs

* fix array of obj not showing up in json editor in test this step

* datatable scales correctly in DisplayResult and scrolling is much more usable

* avoid next button disapearing and changing layout / hurting ux

* nits

* fix bug when renaming module A to B then module C to A, C takes the schema of A

* fix bug with comments in sql repl

* fix aggrid theme randomly not loading

* bindable script

* better delete button in db manager

* property select doesnt exist

* fix all warnings

* delete $flowStateStore[id] on delete

* feat(cli): generate cursor rules on init (#6270)

* create cursor rules on init

* change gen

* add missing resource-type command

* add resource type command in guidance

* add schema option

* revert

* nit

* nit

* add flow guidance

* nit

* chore(main): release 1.515.0 (#6292)

* chore(main): release 1.515.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>

* fix: improved logs for script

* nits logs

* chore(main): release 1.515.1 (#6295)

* chore(main): release 1.515.1

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>

* merge

* even more indexer tracings

* add more tracing logs

* feat: prevent too large results (>500Mb) from OOMing database

* nit naming

* feat: add CA certificate update at startup via environment variable (#6280)

* feat: add CA certificate update at startup via environment variable

Add support for running 'update-ca-certificates' at binary startup
when RUN_UPDATE_CA_CERTIFICATE_AT_START environment variable is set to "true".

- Check for RUN_UPDATE_CA_CERTIFICATE_AT_START env var on startup
- Execute update-ca-certificates command if env var is set to "true"
- Log success/failure appropriately with tracing
- Continue startup even if CA certificate update fails
- Non-blocking implementation with proper error handling

Fixes #6279

Co-authored-by: Ruben Fiszel <rubenfiszel@users.noreply.github.com>

* refactor: extract CA certificate update logic into separate function

Extract the CA certificate update logic from windmill_main() into a
dedicated update_ca_certificates_if_requested() function for better
code organization and maintainability.

Co-authored-by: Ruben Fiszel <rubenfiszel@users.noreply.github.com>

* improvements

---------

Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com>
Co-authored-by: Ruben Fiszel <rubenfiszel@users.noreply.github.com>
Co-authored-by: Alexander Petric <alpetric@users.noreply.github.com>
Co-authored-by: Alexander Petric <alex@windmill.dev>

* fix: indexer collection of job logs before indexing (#6300)

* Add flume as dependecy for indexer

* Update ee-repo-ref

* Remove flags from cargo.toml

* Update ee-repo-ref

* Update ee-repo-ref

* fix rust sdk build error (#6305)

Signed-off-by: pyranota <pyra@duck.com>

* fix broken audit logs filter (#6304)

* rename to from to

* goto fix

* default to false if field not present operator settings (#6301)

* git sync UI improvements (#6303)

* ui improvements round 1

* modal cleanup

* init

* UI refactor

* UI cleanup + refactor

* legacy cleanup

* success model -> github actions, non-ee warnings

* sqlx

* npm check

* ee warning everywhere

* last comments

* formatting

* no hardcoded theme

* claude review improvemenets

* fix: no process relative imports for scripts with codebase

* fix: sqs oidc authentication disconnect #6307

* handle metadata for new scripts happen after commit

* handle_deployment_metadata in a task

* nits

* chore: add windmill-utils-internal package (#6299)

* add utils package

* naming

* cleaning

* add docs

* remove log

* use autogenerated types

* remove old

* fix

* cleaning

* add docs

* chore(main): release 1.516.0 (#6298)

* chore(main): release 1.516.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>

* merge

* indexer improvements

* upgrade tantivy to 0.24.2

* use tantivy fork

* nit warnings

* fix oss build

* improve indexer

* chore: use windmill-utils-internal for cli (#6297)

* add utils package

* naming

* cleaning

* simplify assignPath

* rename old files

* same for locks

* create on confirm

* default true

* use replaceinlinescripts from utils

* use extractscriptfromflows

* make it compile

* cleaning

* use argsigtojson

* fix

* fix missing await

* cleaner

* cleaning

* cleaning

* use in frontend

* add docs

* testing

* remove log

* use autogenerated types

* remove old

* fix

* cleaning

* adapt usage

* draft

* better build script

* fix build

* revert to default creation

* add docs

* remove and rename

* make everything work

* add await

* only if not installed

* add vs code setting

* add to publish action

* fix bc

* safer use of sep

* fix

* do not rename on push

* no publish on release

* use published package on frontend

* nit

* Add dependencies to run sqlx prepare to nix flake (#6309)

* feat(cli): wmill-lock.yaml v2 for easier git merge diffs

* merge

* merge

* all

* all

* rm warnings

* fix styling on aichatinput (#6312)

* fix: use with_capacity back presusre for tantivy directory multipart writes (#6313)

* use with capacity for tantivy directory multi part uploads

* Update ee repo ref

* Update ee-repo-ref

* Update ee-repo-ref

* chore(main): release 1.517.0 (#6310)

* chore(main): release 1.517.0

* Apply automatic changes

---------

Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>

* fix typo on cli build (#6314)

* cleanup

* feat(utils): add flow.yaml validation function (#6316)

* add validateflow function

* cleaner code

* preprocess json

* cleaning

* create specific package

* cleaning

* add tests

* fix: cleanup concurrency_counter automatically + remove orphans keys automatically

* fix: add disabled support to resource picker in schema forms

* fix: add wm_labels to tracing spans

* all

* merge

* all

* fix: delete empty git connection (#6318)

* fix checks

* bun handling

* all

* all?

* all

* all

* update

* all

* update

* check

* fix history

* all

* all

* all

* Remove leftover debug tracing statements

- Remove commented debug trace in jobs.rs for stream output
- Remove commented debug trace in result_stream.rs for stream processing

Co-authored-by: Ruben Fiszel <rubenfiszel@users.noreply.github.com>

* fix test

* all

* handle iter

* fix

---------

Signed-off-by: pyranota <pyra@duck.com>
Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com>
Co-authored-by: Ruben Fiszel <rubenfiszel@users.noreply.github.com>
Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
Co-authored-by: centdix <40307056+centdix@users.noreply.github.com>
Co-authored-by: dieriba <dieriba.pro@gmail.com>
Co-authored-by: rubenfiszel <275584+rubenfiszel@users.noreply.github.com>
Co-authored-by: wendrul <53628737+wendrul@users.noreply.github.com>
Co-authored-by: Diego Imbert <70353967+diegoimbert@users.noreply.github.com>
Co-authored-by: Alexander Petric <alpetric@users.noreply.github.com>
Co-authored-by: Alexander Petric <alex@windmill.dev>
Co-authored-by: pyranota <92104930+pyranota@users.noreply.github.com>
2025-08-06 22:40:42 +00:00

364 lines
11 KiB
Rust

use std::{collections::HashMap, process::Stdio};
use itertools::Itertools;
use serde_json::value::RawValue;
use tokio::{fs::File, io::AsyncWriteExt, process::Command};
use windmill_common::{
error::Error,
worker::{write_file, Connection},
};
use windmill_parser::Arg;
use windmill_parser_nu::parse_nu_signature;
use windmill_queue::{append_logs, CanceledBy, MiniPulledJob};
use crate::{
common::{
create_args_and_out_file, get_reserved_variables, read_result, start_child_process,
OccupancyMetrics,
},
handle_child, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PATH_ENV, PROXY_ENVS,
};
use windmill_common::client::AuthedClient;
const NSJAIL_CONFIG_RUN_NU_CONTENT: &str = include_str!("../nsjail/run.nu.config.proto");
lazy_static::lazy_static! {
static ref NU_PATH: String = std::env::var("NU_PATH").unwrap_or_else(|_| "/usr/bin/nu".to_string());
// TODO(v1):
// static ref PLUGIN_USE_RE: Regex = Regex::new(r#"(?:plugin use )(?<plugin>.*)"#).unwrap();
}
// TODO: Can be generalized and used for other handlers
#[allow(dead_code)]
pub(crate) struct JobHandlerInput<'a> {
pub base_internal_url: &'a str,
pub canceled_by: &'a mut Option<CanceledBy>,
pub client: &'a AuthedClient,
pub parent_runnable_path: Option<String>,
pub conn: &'a Connection,
pub envs: HashMap<String, String>,
pub inner_content: &'a str,
pub job: &'a MiniPulledJob,
pub job_dir: &'a str,
pub mem_peak: &'a mut i32,
pub occupancy_metrics: &'a mut OccupancyMetrics,
pub requirements_o: Option<&'a String>,
pub shared_mount: &'a str,
pub worker_name: &'a str,
}
pub async fn handle_nu_job<'a>(mut args: JobHandlerInput<'a>) -> Result<Box<RawValue>, Error> {
// TODO(v1):
// --- Handle plugins ---
// let plugins = get_plugins(&mut args).await?;
// TODO(v1):
// --- Handle imports ---
// TODO(v1):
// --- Handle relative ---
// --- Wrap and write to fs ---
{
create_args_and_out_file(&args.client, args.job, args.job_dir, args.conn).await?;
File::create(format!("{}/main.nu", args.job_dir))
.await?
.write_all(&wrap(args.inner_content)?.into_bytes())
.await?;
}
// --- Execute ---
{
run(&mut args).await?;
}
// --- Retrieve results ---
{
read_result(&args.job_dir, None).await
}
}
// async fn get_plugins<'a>(
// JobHandlerInput {
// occupancy_metrics,
// mem_peak,
// canceled_by,
// worker_name,
// job,
// db,
// inner_content,
// ..
// }: &mut JobHandlerInput<'a>,
// ) -> Result<Vec<&'a str>, Error> {
// let plugins_dir = concatcp!(NU_CACHE_DIR, "/plugins");
// let nu_version = from_utf8_mut(
// Command::new(NU_PATH.as_str())
// .arg("--version")
// .output()
// .await?
// .stdout
// .as_mut_slice(),
// )
// .map_err(|e| windmill_common::error::Error::ExecutionErr(e.to_string()))?
// .to_owned();
// let plugins = parse_plugin_use(inner_content);
// for plugin in &plugins {
// let mut run_cmd = Command::new(CARGO_PATH.as_str());
// // cargo install nu_plugin_query --version (nu --version); plugin add ~/.cargo/bin/nu_plugin_query
// run_cmd
// // TODO: make it work with env_clear
// // .env_clear()
// .args(&[
// "install",
// "--root",
// plugins_dir,
// "--locked",
// &format!("nu_plugin_{plugin}"),
// "--version",
// &nu_version,
// ])
// .stdout(Stdio::piped())
// .stderr(Stdio::piped());
// #[cfg(windows)]
// nsjail_cmd.env("SystemRoot", SYSTEM_ROOT.as_str());
// let child = start_child_process(run_cmd, "cargo").await?;
// // handle_child::handle_child(
// // &job.id,
// // db,
// // mem_peak,
// // canceled_by,
// // child,
// // !*DISABLE_NSJAIL,
// // worker_name,
// // &job.workspace_id,
// // "cargo",
// // job.timeout,
// // false,
// // &mut Some(occupancy_metrics),
// // )
// // .await?;
// }
// Ok(plugins)
// }
// fn parse_plugin_use(inner_content: &str) -> Vec<&str> {
// let mut plugins = vec![];
// // TODO: Ignore plugins with # in the beginning
// for cap in PLUGIN_USE_RE.captures_iter(inner_content).into_iter() {
// if let Some(mat) = cap.name("plugin") {
// plugins.push(mat.as_str());
// }
// }
// plugins
// }
/// Wraps content script
/// that upon execution reads args.json (which are piped and transformed from previous flow step or top level inputs)
/// Also wrapper takes output of program and serializes to result.json (Which windmill will know how to use later)
fn wrap(inner_content: &str) -> Result<String, Error> {
let sig = parse_nu_signature(inner_content)?;
let spread = sig
.args
.clone()
.into_iter()
.map(|Arg { name, typ, has_default, .. }| {
// Apply additional input transformation
let transformation = format!(
"| if $in != null {{ {} }} else {{ $in }}",
match typ {
// JSON converts X.0 to X and nu can't coerce type automatically
windmill_parser::Typ::Datetime => "into datetime",
windmill_parser::Typ::Bytes => "into binary",
windmill_parser::Typ::Float => "into float",
// Ident
_ => "$in",
}
);
let nullguard = if has_default || matches!(typ, windmill_parser::Typ::Unknown) {
"".to_owned()
} else {
format!("| nullguard {name}")
};
format!("\n\t\t\t($parsed_args.{name}? {nullguard} {transformation}) ",)
})
.collect_vec()
.join(" ");
Ok(
r#"
$env.config.table.mode = 'basic'
def nullguard [ name: string ] {
if ($in == null) {
panic $"argument `($name)` of main function can't be null"
}
$in
}
# TODO: Probably needs rework in order for LSP to work
def get_variable [ pat ] {
let addr = $"($env.BASE_INTERNAL_URL)/api/w/($env.WM_WORKSPACE)/variables/get_value/($pat)" ;
http get -H ["Authorization", $"Bearer ($env.WM_TOKEN)"] $addr | return $in
}
def get_resource [ pat ] {
let addr = $"($env.BASE_INTERNAL_URL)/api/w/($env.WM_WORKSPACE)/resources/get_value_interpolated/($pat)" ;
http get -H ["Authorization", $"Bearer ($env.WM_TOKEN)"] $addr | return $in
}
def 'main --wrapped' [] {
let parsed_args = open args.json
(main SPREAD
) | to json | save -f result.json
}
INNER_CONTENT
"#
.replace("INNER_CONTENT", inner_content)
.replace("SPREAD", &spread), // .replace("TRANSFORM", transform)
)
}
async fn run<'a>(
JobHandlerInput {
occupancy_metrics,
mem_peak,
canceled_by,
worker_name,
job,
conn,
job_dir,
shared_mount,
client,
parent_runnable_path,
envs,
base_internal_url,
..
}: &mut JobHandlerInput<'a>,
// plugins: Vec<&'a str>,
) -> Result<(), Error> {
let reserved_variables =
get_reserved_variables(job, &client.token, conn, parent_runnable_path.clone()).await?;
let child = if !cfg!(windows) && !*DISABLE_NSJAIL {
append_logs(
&job.id,
&job.workspace_id,
format!("\n\n--- ISOLATED NU CODE EXECUTION ---\n"),
conn,
)
.await;
write_file(
job_dir,
"run.config.proto",
&NSJAIL_CONFIG_RUN_NU_CONTENT
.replace("{JOB_DIR}", job_dir)
.replace("{NU_PATH}", &NU_PATH)
.replace("{SHARED_MOUNT}", &shared_mount)
.replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()),
)?;
let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str());
nsjail_cmd
.env_clear()
.current_dir(job_dir)
.env("PATH", PATH_ENV.as_str())
.env("BASE_INTERNAL_URL", base_internal_url)
.envs(envs)
.envs(reserved_variables)
.envs(PROXY_ENVS.clone())
.args(vec![
"--config",
"run.config.proto",
"--",
NU_PATH.as_str(),
"/tmp/main.nu",
"--wrapped",
])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
start_child_process(nsjail_cmd, NSJAIL_PATH.as_str()).await?
} else {
append_logs(
&job.id,
&job.workspace_id,
format!("\n\n--- NU CODE EXECUTION ---\n"),
&conn,
)
.await;
// let plugin_registry = format!("{job_dir}/plugin-registry");
// File::create(&plugin_registry).await?;
//
let mut cmd = Command::new(if cfg!(windows) {
"nu"
} else {
NU_PATH.as_str()
});
cmd.env_clear()
.current_dir(job_dir.to_owned())
.env("PATH", PATH_ENV.as_str())
.env("BASE_INTERNAL_URL", base_internal_url)
.envs(envs)
.envs(reserved_variables)
.envs(PROXY_ENVS.clone())
.args(&[
"main.nu",
"--wrapped",
// TODO(v1):
// "--plugins",
// &format!(
// "[{}]",
// plugins
// .into_iter()
// .map(|pl| format!("{NU_CACHE_DIR}/plugins/bin/nu_plugin_{pl}"))
// .collect_vec()
// .join(",")
// ),
])
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(windows)]
{
cmd.env("SystemRoot", crate::SYSTEM_ROOT.as_str())
.env("USERPROFILE", crate::USERPROFILE_ENV.as_str())
.env(
"TMP",
std::env::var("TMP").unwrap_or_else(|_| String::from("/tmp")),
);
}
start_child_process(cmd, "nu").await?
};
handle_child::handle_child(
&job.id,
conn,
mem_peak,
canceled_by,
child,
!*DISABLE_NSJAIL,
worker_name,
&job.workspace_id,
"nu",
job.timeout,
false,
&mut Some(occupancy_metrics),
None,
)
.await?;
Ok(())
}
// #[cfg(test)]
// mod test {
// use super::parse_plugin_use;
// #[test]
// fn test_nu_plugin_use() {
// let content = r#"
// plugin use foo
// plugin use bar
// plugin use baz
// plugin use meh
// "#;
// assert_eq!(
// vec!["foo", "bar", "baz", "meh"], //
// parse_plugin_use(content)
// );
// }
// }