diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..eab2707 --- /dev/null +++ b/Makefile @@ -0,0 +1,9 @@ +.PHONY: all +all: priv/xdg-open + # noop + +.PHONY: priv/xdg-open +priv/xdg-open: + mkdir -p priv + curl https://raw.githubusercontent.com/sindresorhus/open/refs/heads/main/xdg-open > priv/xdg-open + chmod +x priv/xdg-open diff --git a/README.md b/README.md index 33a5b66..a22057d 100644 --- a/README.md +++ b/README.md @@ -36,7 +36,7 @@ import child_process import child_process/stdio // Run vim interactively - stdio is inherited so you can use it normally! -child_process.new("vim") +child_process.from_name("vim") |> child_process.arg("README.md") |> child_process.stdio(stdio.inherit()) |> child_process.run() @@ -49,7 +49,7 @@ import child_process import child_process/stdio // Watch mode with live output -child_process.new("tailwindcss") +child_process.from_name("tailwindcss") |> child_process.args(["-i", "input.css", "-o", "output.css", "--watch"]) |> child_process.stdio(stdio.lines(fn(line) { io.println(line) diff --git a/dev/open.gleam b/dev/open.gleam new file mode 100644 index 0000000..4595f3a --- /dev/null +++ b/dev/open.gleam @@ -0,0 +1,6 @@ +import child_process + +pub fn main() { + echo child_process.open("https://gleam.run") + // echo child_process.reveal("./gleam.toml") +} diff --git a/dev/vi.gleam b/dev/vi.gleam index 8fd11cd..0eb5886 100644 --- a/dev/vi.gleam +++ b/dev/vi.gleam @@ -2,8 +2,7 @@ import child_process import child_process/stdio pub fn main() { - child_process.new_with_path("hx") - |> child_process.stdio(stdio.inherit()) - |> child_process.run + child_process.from_name("hx") + |> child_process.run(stdio.inherit()) |> echo } diff --git a/src/child_process.gleam b/src/child_process.gleam index 2ca67c2..9b65739 100644 --- a/src/child_process.gleam +++ b/src/child_process.gleam @@ -1,5 +1,5 @@ import child_process/internal.{Builder, Linux, Macos, Windows} -import child_process/stdio.{type Stdio} +import child_process/stdio.{type Mode, type Stdio} import gleam/dict import gleam/list import gleam/option.{type Option} @@ -7,7 +7,9 @@ import gleam/result import gleam/string @target(erlang) -import gleam/erlang/process +import gleam/erlang/port.{type Port} +@target(erlang) +import gleam/erlang/process.{type Pid} @target(erlang) import gleam/otp/actor @target(erlang) @@ -15,69 +17,21 @@ import gleam/otp/factory_supervisor @target(erlang) import gleam/otp/supervision -// TODO(v2): steal louis error type -// TODO(v2): support the stdio folds in the synchronous case? -// TODO(v2) rename new with path as its confusing -// TODO: figure out if we always have to spawn the intermediate process -// TODO: null, inherit, stream, lines -// if on_error is optional then null, inherit and stream -// louis has custom `select` exposed which is also a callback so i guess its just about the process -// not using that leaks messages -// if we don't pass use_stdio/eof/exit_status i think it doesnt send messages? -// -// select_stream -// select_exit_status -// ok but to have a new selector you would have to return that right and keep the old one around etc -// a subject enables us to go the other way: have a configured selector and add new ways of receiving messages on the fly -// basically we have messages but we have to transform them on receive. selectors allow us that but subjects don't -// ... gleam_otp is bad??? really??? -// ... can they be the same type? -// -// select cannot be a builder though, or it would be a bit awkward since we need to return that selector -// subject(Opaque) that you pass in, select that handles these? -// -// -// would have the same interface with callbacks, but additionally a selector -// if only the types where the other way around :(( -// -// -// -// that leaves lines that is special -// select_lines is impossible in a way as it requires state -// it would need to follow a different interface, like fn(data, is_full_line) -// -// -// -// i suppose the way would be to mirror the selector api?? -// still the otp design splits configuration and usage, so we will always have a resource problem no -// eh we kinda have the same project if you never selected from the subject ever its not a channel -// -// -// so we can have our callbacks and then use `select_spawn` to skip the process -// or maybe we can skip the old version then? probably for everything except for lines -// -// // TODO: issue for "add unqualified import" for type constructors in cases /// A builder for configuring and spawning child processes. /// -/// Use `new()` or `new_with_path()` to create a builder, then chain +/// Use `from_file()` or `from_name()` to create a builder, then chain /// configuration methods like `arg()`, `env()`, and `cwd()` before /// calling `run()` or `spawn()`. /// /// ## Example /// ```gleam -/// child_process.new("node") +/// child_process.from_name("node") /// |> child_process.arg("server.js") /// |> child_process.env("PORT", "3000") /// |> child_process.cwd("/app") -/// |> child_process.stdio(stdio.lines(fn(line) { -/// io.println(line) -/// })) -/// |> child_process.on_exit(fn(code) { -/// io.println("Server exited with code: " <> int.to_string(code)) -/// }) -/// |> child_process.spawn() +/// |> child_process.run(stdio.capture(True)) /// ``` pub type Builder = internal.Builder @@ -100,20 +54,58 @@ pub type Output { /// - `-2` indicates the process failed to start (only on JavaScript target) /// - Other non-zero values are process-specific error codes /// - /// **Output capture:** The output field will be empty if `stdio.inherit()` - /// was used, since output goes directly to the terminal. + /// When no output was configured to be captured, `output` will be empty. Output(status_code: Int, output: String) } /// An error that occurred when trying to start a process. pub type StartError { + /// The operating system does not have enough memory to spawn this process. + OutOfMemory + /// The operating system process table is full. + OsProcessLimitReached + /// The operating system file descriptor table is full. + OsFileLimitReached + /// The command is too long. + CommandTooLong + /// The runtime process has too many open file descriptors. + NotEnoughFileDescriptors /// The executable file was not found at the specified path. FileNotFound(String) /// The file exists but is not executable. FileNotExecutable(String) - /// A system resource limit prevented the process from starting. - /// This could be an operating system limit or a VM limit. - SystemLimit + /// A runtime resource limit prevented the process from starting. + RuntimeLimitReached +} + +/// Convert a `StartError` into a human-readable string. +pub fn describe_start_error(error: StartError) { + case error { + OutOfMemory -> "Not enough memory" + OsProcessLimitReached -> "OS process limit reached" + OsFileLimitReached -> "OS file descriptor limit reached" + CommandTooLong -> "Command too long" + NotEnoughFileDescriptors -> "Runtime has too many open file descriptors" + FileNotFound(executable) -> "File not found: " <> executable + FileNotExecutable(executable) -> "File not executable: " <> executable + RuntimeLimitReached -> "Runtime limit reached" + } +} + +/// An error that occurred while sending data to the process. +pub type WriteError { + /// The process can't be written to right now or writing was cancelled. + WriteAborted + /// The process you tried to write to already exited. + ProcessExited +} + +/// Convert a `WriteError` into a human-readable string. +pub fn describe_write_error(error: WriteError) { + case error { + WriteAborted -> "Write aborted, try again" + ProcessExited -> "The process already exited" + } } // -- BUILDER ----------------------------------------------------------------- @@ -122,11 +114,11 @@ pub type StartError { /// /// ## Example /// ```gleam -/// new("/usr/bin/git") +/// from_file("/usr/bin/git") /// |> arg("status") -/// |> run() +/// |> run(stdio.capture(True)) /// ``` -pub fn new(executable_path: String) -> Builder { +pub fn from_file(executable_path: String) -> Builder { new_with_executable(internal.File(executable_path)) } @@ -134,11 +126,11 @@ pub fn new(executable_path: String) -> Builder { /// /// ## Example /// ```gleam -/// new_with_path("git") +/// from_name("git") /// |> arg("status") -/// |> run() +/// |> run(stdio.capture(True)) /// ``` -pub fn new_with_path(executable_name: String) -> Builder { +pub fn from_name(executable_name: String) -> Builder { new_with_executable(internal.Named(executable_name)) } @@ -148,8 +140,6 @@ fn new_with_executable(executable: internal.Executable) -> Builder { reverse_args: [], environment: dict.new(), working_directory: option.None, - stdio: stdio.null(), - on_exit: fn(_) { Nil }, ) } @@ -190,48 +180,28 @@ pub fn cwd(builder: Builder, directory: String) -> Builder { Builder(..builder, working_directory: option.Some(directory)) } -/// Configure stdio handling for the process. -pub fn stdio(builder: Builder, stdio: Stdio) -> Builder { - Builder(..builder, stdio:) -} - -/// Configure a callback that is called when the process exits. -/// -/// ## Example -/// ```gleam -/// new("vite") -/// |> on_exit(fn(status_code) { -/// process.end(self, ViteExited(status_code)) -/// }) -/// ``` -pub fn on_exit(builder: Builder, on_exit callback: fn(Int) -> Nil) -> Builder { - Builder(..builder, on_exit: callback) -} - // -- STARTING ---------------------------------------------------------------- /// Run the process synchronously and wait for it to complete. /// -/// If the stdio is not explicitly set to `inherit`, the stdio configuration -/// will be replaced and all output will be captured. +/// When mode is `capture(_)`, output is captured and returned in `Output.output`. +/// When mode is `null()` or `inherit()`, output is empty. /// -/// Returns the exit status and captured output. The `on_exit` -/// callback will be called before this function returns. -pub fn run(builder: Builder) -> Result(Output, StartError) { - let Builder( - executable:, - reverse_args:, - environment:, - working_directory: cwd, - stdio:, - on_exit:, - ) = builder +/// ## Example +/// ```gleam +/// child_process.from_file("/bin/echo") +/// |> child_process.arg("hello") +/// |> child_process.run(stdio.capture(True)) +/// ``` +pub fn run(builder: Builder, mode mode: Mode) -> Result(Output, StartError) { + let Builder(executable:, reverse_args:, environment:, working_directory: cwd) = + builder use executable <- result.try(resolve_executable(executable)) let args = list.reverse(reverse_args) let env = dict.to_list(environment) - do_run(executable, args, env, cwd, stdio, on_exit) + do_run(executable, args, env, cwd, mode) } @external(erlang, "child_process_ffi", "run") @@ -241,31 +211,76 @@ fn do_run( args: List(String), env: List(#(String, String)), working_directory: Option(String), - stdio: Stdio, - on_exit: fn(Int) -> Nil, + mode: Mode, ) -> Result(Output, StartError) -/// Spawn the process asynchronously, returning a handle to interact with it. +/// Spawn the process asynchronously with a stdio handler. /// -/// > **Note:** Node delivers errors asynchronously, so this implementation -/// > makes various checks beforehand to be reasonably sure that the call to -/// > spawn will succeed. If the process fails to start, `on_exit` will report -/// > a status_code of -2. -pub fn spawn(builder: Builder) -> Result(Process, StartError) { - let Builder( - executable:, - reverse_args:, - environment:, - working_directory: cwd, - stdio:, - on_exit:, - ) = builder +/// The stdio configuration determines both the I/O mode (how the port is +/// configured) and the callbacks (how data is consumed). The handler is +/// applied immediately after spawning. +/// +/// On Erlang, this spawns an extra process that manages the port. +/// +/// +/// ## Example +/// ```gleam +/// let assert Ok(proc) = +/// child_process.from_name("node") +/// |> child_process.arg("server.js") +/// |> child_process.spawn(stdio.lines(fn(line) { io.println(line) })) +/// ``` +/// +/// > **Note:** On JavaScript, Node delivers errors asynchronously, so this +/// > implementation makes various checks beforehand to be reasonably sure +/// > that the call will succeed. +pub fn spawn( + builder: Builder, + stdio stdio: Stdio, +) -> Result(Process, StartError) { + let Builder(executable:, reverse_args:, environment:, working_directory: cwd) = + builder use executable <- result.try(resolve_executable(executable)) let args = list.reverse(reverse_args) let env = dict.to_list(environment) - do_spawn(executable, args, env, cwd, stdio, on_exit) + do_spawn(executable, args, env, cwd, stdio) +} + +/// Spawn the process asynchronously without attaching +/// an extra process or event listeners. +/// +/// On Erlang, the port is opened in the calling process. Port messages +/// arrive in the caller's mailbox and can be received via `stdio.select()`. +/// +/// ## Example +/// ```gleam +/// let assert Ok(proc) = +/// child_process.from_file("my_program") +/// |> child_process.spawn_raw(stdio.capture(True)) +/// +/// let selector = +/// process.new_selector() +/// |> stdio.select(proc, HandleData, HandleExit) +/// ``` +/// +/// > **Note:** On JavaScript, Node delivers errors asynchronously, so this +/// > implementation makes various checks beforehand to be reasonably sure +/// > that the call will succeed. +/// +pub fn spawn_raw( + builder: Builder, + mode mode: Mode, +) -> Result(Process, StartError) { + let Builder(executable:, reverse_args:, environment:, working_directory: cwd) = + builder + + use executable <- result.try(resolve_executable(executable)) + let args = list.reverse(reverse_args) + let env = dict.to_list(environment) + + do_spawn_raw(executable, args, env, cwd, mode) } @external(erlang, "child_process_ffi", "spawn") @@ -276,7 +291,16 @@ fn do_spawn( env: List(#(String, String)), working_directory: Option(String), stdio: Stdio, - on_exit: fn(Int) -> Nil, +) -> Result(Process, StartError) + +@external(erlang, "child_process_ffi", "spawn_raw") +@external(javascript, "./child_process_ffi.mjs", "spawn_raw") +fn do_spawn_raw( + executable: String, + args: List(String), + env: List(#(String, String)), + working_directory: Option(String), + mode: Mode, ) -> Result(Process, StartError) /// Execute a command with arguments in a specific directory. @@ -290,10 +314,10 @@ pub fn exec( with arguments: List(String), in directory: String, ) -> Result(Output, StartError) { - new(command) + from_name(command) |> args(arguments) |> cwd(directory) - |> run + |> run(stdio.capture(True)) } /// Execute a shell command string using the system shell. @@ -306,7 +330,10 @@ pub fn exec( pub fn shell(cmd: String) -> Result(String, StartError) { let #(shell_exe, shell_args) = shell_command(cmd) - let result = new(shell_exe) |> args(shell_args) |> run + let result = + from_file(shell_exe) + |> args(shell_args) + |> run(stdio.capture(True)) case result { Ok(Output(status_code: _, output:)) -> Ok(output) @@ -329,12 +356,20 @@ fn resolve_executable( // -- ERLANG-SPECIFICS ------------------------------------------------------- +@target(erlang) +/// Recover the Erlang port from a process. +@external(erlang, "gleam@function", "identity") +pub fn port(process: Process) -> Port + @target(erlang) /// Create a child specification for use with OTP supervisors. /// -/// Using the `Permanent` restart strategy, the process will always be -/// restarted if it exits (including normal exits). This allows you to add -/// long running external processes like servers or daemons to your OTP app. +/// The stdio configuration determines both how I/O is handled and how +/// data is consumed. The intermediate process created by the handler +/// is what the supervisor monitors. +/// +/// When the supervisor shuts down, the intermediate process receives a +/// shutdown signal and sends SIGTERM to the OS process. /// /// ## Example /// ```gleam @@ -344,26 +379,31 @@ fn resolve_executable( /// |> supervisor.add( /// child_process.new("node") /// |> child_process.arg("server.js") -/// |> child_process.supervised() +/// |> child_process.supervised(stdio.lines(fn(line) { +/// io.println(line) +/// })) /// ) /// ``` -pub fn supervised(builder: Builder) -> supervision.ChildSpecification(Process) { +pub fn supervised( + builder: Builder, + stdio stdio: Stdio, +) -> supervision.ChildSpecification(Process) { use <- supervision.worker - spawn_child(builder) + spawn_child(builder, stdio) } @target(erlang) /// Create a factory supervisor builder for dynamically spawning supervised processes. /// -/// Keep in mind that a non-zero exit code still indicates success and the Erlang -/// process will exit normally. +/// The stdio configuration determines how each child's I/O is handled. /// /// ## Example /// ```gleam /// import gleam/otp/factory_supervisor /// /// let assert Ok(factory) = -/// child_process.factory() |> factory_supervisor.start() +/// child_process.factory(stdio: stdio.ignore()) +/// |> factory_supervisor.start() /// /// let worker_config = /// child_process.new("node") @@ -371,27 +411,28 @@ pub fn supervised(builder: Builder) -> supervision.ChildSpecification(Process) { /// /// factory_supervisor.start_child(factory, worker_config) /// ``` -pub fn factory() -> factory_supervisor.Builder(Builder, Process) { +pub fn factory( + stdio stdio: Stdio, +) -> factory_supervisor.Builder(Builder, Process) { use builder: Builder <- factory_supervisor.worker_child - spawn_child(builder) + spawn_child(builder, stdio) } @target(erlang) -fn spawn_child(builder: Builder) { - case spawn(builder) { - Error(FileNotFound(path)) -> - Error(actor.InitFailed("File not found: " <> path)) - Error(FileNotExecutable(path)) -> - Error(actor.InitFailed("File not executable: " <> path)) - Error(SystemLimit) -> Error(actor.InitFailed("system limit")) - - Ok(process) -> Ok(actor.Started(erlang_process_pid(process), process)) +fn spawn_child(builder: Builder, stdio: Stdio) { + case spawn(builder, stdio:) { + Error(error) -> Error(actor.InitFailed(describe_start_error(error))) + Ok(process) -> + case port_owner(process) { + Ok(pid) -> Ok(actor.Started(pid, process)) + Error(_) -> Error(actor.InitFailed("Process exited")) + } } } @target(erlang) -@external(erlang, "child_process_ffi", "erlang_process_pid") -fn erlang_process_pid(process: Process) -> process.Pid +@external(erlang, "child_process_ffi", "port_owner") +fn port_owner(process: Process) -> Result(Pid, Nil) // -- PROCESS API ------------------------------------------------------------- @@ -401,10 +442,10 @@ fn erlang_process_pid(process: Process) -> process.Pid /// as the stdin stream is connected directly to the parent's terminal. @external(erlang, "child_process_ffi", "write") @external(javascript, "./child_process_ffi.mjs", "write") -pub fn write(process: Process, text: String) -> Nil +pub fn write(process: Process, text: String) -> Result(Nil, WriteError) /// Write a line of text to the process's stdin. -pub fn writeln(process: Process, text: String) -> Nil { +pub fn writeln(process: Process, text: String) -> Result(Nil, WriteError) { write(process, text <> newline()) } @@ -433,7 +474,12 @@ pub fn kill(process: Process) -> Nil /// Get the operating system process ID. @external(erlang, "child_process_ffi", "os_process_id") @external(javascript, "./child_process_ffi.mjs", "os_process_id") -pub fn os_process_id(process: Process) -> Int +pub fn os_process_id(process: Process) -> Result(Int, Nil) + +/// Returns `True` if the process is still running. +pub fn is_running(process: Process) { + result.is_ok(os_process_id(process)) +} // -- UTILS ------------------------------------------------------------------- @@ -478,24 +524,24 @@ fn shell_command(command_line: String) -> #(String, List(String)) { pub fn open(target: String) -> Result(Nil, StartError) { let builder = case internal.os() { Windows -> - new_with_path("cmd.exe") + from_name("cmd.exe") |> arg("/c") |> arg("start") |> arg(target) Macos -> - new_with_path("open") + from_name("open") |> arg(target) Linux -> { let priv = priv_directory() - new(priv <> "/background") + from_file(priv <> "/background") |> arg(priv <> "/xdg-open") |> arg(target) } } - case run(builder |> stdio(stdio.null())) { + case run(builder, stdio.null()) { Ok(Output(0, _)) -> Ok(Nil) Ok(Output(_, _)) -> Error(FileNotExecutable(target)) Error(error) -> Error(error) @@ -513,16 +559,16 @@ pub fn reveal(path: String) -> Result(Nil, StartError) { let os = internal.os() let builder = case os { Windows -> - new_with_path("explorer.exe") + from_name("explorer.exe") |> arg("/select,\"" <> string.replace(path, "/", "\\") <> "\"") Macos -> - new_with_path("open") + from_name("open") |> arg("--reveal") |> arg(path) Linux -> - new_with_path("gdbus") + from_name("gdbus") |> arg("call") |> arg("--session") |> arg("--dest") @@ -535,7 +581,7 @@ pub fn reveal(path: String) -> Result(Nil, StartError) { |> arg("") } - case run(builder |> stdio(stdio.null())) { + case run(builder, stdio.null()) { // Explorer returns exit code 1 even on success Ok(Output(0, _)) -> Ok(Nil) Ok(Output(1, _)) if os == Windows -> Ok(Nil) diff --git a/src/child_process/internal.gleam b/src/child_process/internal.gleam index 2333833..4dbb2e0 100644 --- a/src/child_process/internal.gleam +++ b/src/child_process/internal.gleam @@ -3,18 +3,18 @@ import gleam/dict.{type Dict} import gleam/option.{type Option} -/// Type-erased value for storing arbitrary types +/// Type-erased value for storing arbitrary types. pub type Any -/// Convert any value to type-erased Any +/// Convert any value to type-erased Any. @external(erlang, "gleam@function", "identity") @external(javascript, "../../gleam_stdlib/gleam/function.mjs", "identity") -pub fn to_any(value: a) -> Any +fn to_any(value: a) -> Any -/// Convert type-erased Any back to original type +/// Convert type-erased Any back to original type. @external(erlang, "gleam@function", "identity") @external(javascript, "../../gleam_stdlib/gleam/function.mjs", "identity") -pub fn from_any(value: Any) -> a +fn from_any(value: Any) -> a /// Internal representation of a process builder. pub type Builder { @@ -23,8 +23,6 @@ pub type Builder { reverse_args: List(String), environment: Dict(String, String), working_directory: Option(String), - stdio: Stdio, - on_exit: fn(Int) -> Nil, ) } @@ -34,15 +32,34 @@ pub type Executable { Named(name: String) } -/// Internal representation of stdio configuration with type-erased accumulator -pub type Stdio { +/// Internal representation of I/O mode configuration. +pub type Mode { Null Inherit - Reduce( + Capture(capture_stderr: Bool) +} + +/// Describes how to handle I/O for a spawned process. +pub type Stdio { + Stdio( + mode: Mode, initial: Any, on_data: fn(Any, BitArray) -> Any, - on_exit: fn(Any, Int) -> Any, - capture_stderr: Bool, + on_exit: fn(Any, Int) -> Nil, + ) +} + +pub fn stdio( + mode mode: Mode, + initial initial: state, + on_data on_data: fn(state, BitArray) -> state, + on_exit on_exit: fn(state, Int) -> Nil, +) -> Stdio { + Stdio( + mode:, + initial: to_any(initial), + on_data: fn(state, bits) { to_any(on_data(from_any(state), bits)) }, + on_exit: fn(state, status) { on_exit(from_any(state), status) }, ) } diff --git a/src/child_process/line_buffer.gleam b/src/child_process/line_buffer.gleam new file mode 100644 index 0000000..54daba8 --- /dev/null +++ b/src/child_process/line_buffer.gleam @@ -0,0 +1,65 @@ +import gleam/list +import gleam/option.{type Option, None, Some} +import splitter + +/// A buffer for accumulating partial line data from a port. +/// +/// Use this when processing data from a child process via streams +/// and you want line-by-line processing. +pub opaque type LineBuffer { + LineBuffer(buffer: String, splitter: splitter.Splitter) +} + +/// Create a new line buffer. +pub fn new() -> LineBuffer { + LineBuffer(buffer: "", splitter: splitter.new([newline()])) +} + +/// Feed data into the line buffer. +/// +/// Returns a list of complete lines (including their newline characters) +/// and the updated buffer containing any partial line data. +/// +/// ## Example +/// ```gleam +/// let buffer = line_buffer.new() +/// let #(lines, buffer) = line_buffer.feed(buffer, <<"hello\nworld">>) +/// // lines = ["hello\n"] +/// // buffer still contains "world" +/// ``` +pub fn feed(buffer: LineBuffer, data: BitArray) -> #(List(String), LineBuffer) { + let text = lossy_bits_to_string(data) + let #(lines, remaining) = + split_all_lines(buffer.splitter, buffer.buffer <> text, []) + #(lines, LineBuffer(..buffer, buffer: remaining)) +} + +/// Flush any remaining data in the buffer. +/// +/// Call this when the process exits to get any final partial line +/// that didn't end with a newline. +pub fn flush(buffer: LineBuffer) -> Option(String) { + case buffer.buffer { + "" -> None + text -> Some(text) + } +} + +fn split_all_lines( + line_splitter: splitter.Splitter, + data: String, + acc: List(String), +) -> #(List(String), String) { + case splitter.split_after(line_splitter, data) { + #(line, "") -> #(list.reverse(acc), line) + #(line, rest) -> split_all_lines(line_splitter, rest, [line, ..acc]) + } +} + +@external(erlang, "child_process_ffi", "newline") +@external(javascript, "../child_process_ffi.mjs", "newline") +fn newline() -> String + +@external(erlang, "child_process_ffi", "bits_to_string") +@external(javascript, "../child_process_ffi.mjs", "bits_to_string") +fn lossy_bits_to_string(bits: BitArray) -> String diff --git a/src/child_process/stdio.gleam b/src/child_process/stdio.gleam index d037ad3..68c1612 100644 --- a/src/child_process/stdio.gleam +++ b/src/child_process/stdio.gleam @@ -1,17 +1,19 @@ -import child_process/internal.{Inherit, Null, Reduce, from_any, to_any} -import splitter +import child_process/internal.{Capture, Inherit, Null, Stdio} +import child_process/line_buffer +import gleam/list +import gleam/option.{None, Some} -/// Configuration for how a process's standard I/O streams are handled. -pub type Stdio = - internal.Stdio +@target(erlang) +import child_process/internal.{type Process} as _ +@target(erlang) +import gleam/erlang/process.{type Selector} -/// Get the native newline character for the operating system. -@external(erlang, "child_process_ffi", "newline") -@external(javascript, "../child_process_ffi.mjs", "newline") -fn newline() -> String +/// I/O mode configuration for a child process port. +pub type Mode = + internal.Mode -/// Discard all output from the process. This is the default. -pub fn null() -> Stdio { +/// Discard all output from the process. +pub fn null() -> Mode { Null } @@ -28,149 +30,149 @@ pub fn null() -> Stdio { /// programs. /// /// For example: `erl -noshell -pa build/dev/erlang/*/ebin -run module main` -/// -/// ## Example -/// ```gleam -/// // Launch vim - user can interact with it directly -/// child_process.new("vim") -/// |> child_process.arg("file.txt") -/// |> child_process.stdio(stdio.inherit()) -/// |> child_process.run() -/// ``` -pub fn inherit() -> Stdio { +pub fn inherit() -> Mode { Inherit } -fn reduce( - from state: state, - on_data on_data: fn(state, String) -> state, - on_exit on_exit: fn(state, Int) -> result, -) -> Stdio { - let on_data = fn(state, bits) { on_data(state, lossy_bits_to_string(bits)) } - reduce_bits(state, on_data, on_exit) +/// Capture the process's stdout (and optionally stderr) so it can be read. +/// +/// When used with `run()`, output is automatically collected and returned. +pub fn capture(capture_stderr capture_stderr: Bool) -> Mode { + Capture(capture_stderr:) } -@external(erlang, "child_process_ffi", "bits_to_string") -@external(javascript, "../child_process_ffi.mjs", "bits_to_string") -fn lossy_bits_to_string(bits: BitArray) -> String +// -- STDIO ------------------------------------------------------------------- + +/// Describes how to handle I/O for a spawned process. +/// +/// A `Stdio` bundles an I/O mode (how the port is configured) with +/// callbacks (how data is consumed). Create one with `lines()`, `chunks()`, +/// `collect()`, `exit()`, or `ignore()`. +pub type Stdio = + internal.Stdio -fn reduce_bits( - from state: state, +/// Override whether stderr is captured in addition to stdout for this `Stdio`. +pub fn capture_stderr(stdio: Stdio, capture capture: Bool) -> Stdio { + case stdio.mode { + Capture(_) -> Stdio(..stdio, mode: Capture(capture_stderr: capture)) + _ -> stdio + } +} + +/// Create a custom data handler. Most of the time, you would want to use one of +/// the predefined handlers instead. +pub fn custom( + initial state: state, on_data on_data: fn(state, BitArray) -> state, - on_exit on_exit: fn(state, Int) -> result, + on_exit on_exit: fn(state, Int) -> Nil, ) -> Stdio { - Reduce( - to_any(state), - fn(acc, chunk) { to_any(on_data(from_any(acc), chunk)) }, - fn(acc, code) { to_any(on_exit(from_any(acc), code)) }, - False, - ) + internal.stdio(Capture(capture_stderr: True), state, on_data, on_exit) } -/// Stream raw output data as it arrives. -/// The callback receives chunks of data as Strings. -/// -/// ## Example -/// ```gleam -/// stdio.stream(fn(data) { -/// process.send(self, ProcessSentData(data)) -/// }) -/// ``` -pub fn stream(on_data callback: fn(String) -> Nil) -> Stdio { - reduce(Nil, fn(_state, chunk) { callback(chunk) }, fn(_state, _code) { Nil }) -} - -/// Stream raw binary data as it arrives. -/// The callback receives chunks of data as BitArrays. -/// -/// ## Example -/// ```gleam -/// stdio.stream_bits(fn(data) { -/// process.send(self, ProcessSentData(data)) -/// }) -/// ``` -pub fn stream_bits(on_data callback: fn(BitArray) -> Nil) -> Stdio { - reduce_bits(Nil, fn(_, chunk) { callback(chunk) }, fn(_, _) { Nil }) +/// No-op handler that discards all output. +pub fn ignore() -> Stdio { + internal.stdio(Null, Nil, fn(state, _) { state }, fn(_, _) { Nil }) } /// Stream output line-by-line. /// -/// The callback receives complete lines as Strings, including their newline -/// character. Note that if the last line does not contain a newline character, -/// the received string won't contain one either. +/// The callback receives complete lines as Strings, including their newline character. +/// When the process exits, any remaining buffered data is delivered as a final +/// callback without a trailing newline. /// -/// **Important:** If the process outputs partial lines (data without a newline), -/// that data will be buffered until either a newline arrives or the process exits. -/// When the process exits, any remaining buffered data will be delivered as a -/// final callback, even without a trailing newline. +/// Captures stderr by default. Use `capture_stderr(False)` to override. /// -/// ## Example -/// ```gleam -/// stdio.lines(fn(line) { -/// io.println(line) -/// }) -/// ``` pub fn lines(on_line callback: fn(String) -> Nil) -> Stdio { - let line_splitter = splitter.new([newline()]) - reduce( - from: "", - on_data: fn(buffer, chunk) { - process_lines(line_splitter, buffer <> chunk, callback) + custom( + line_buffer.new(), + on_data: fn(buffer, data) { + let #(lines, buffer) = line_buffer.feed(buffer, data) + list.each(lines, callback) + buffer }, on_exit: fn(buffer, _code) { - case buffer { - "" -> Nil - _ -> callback(buffer) + case line_buffer.flush(buffer) { + Some(remaining) -> callback(remaining) + None -> Nil } }, ) } -fn process_lines( - line_splitter: splitter.Splitter, - data: String, - callback: fn(String) -> Nil, -) -> String { - case splitter.split_after(line_splitter, data) { - #(line, "") -> line - #(line, rest) -> { - callback(line) - process_lines(line_splitter, rest, callback) - } - } +/// Stream output data as it arrives. The callback receives chunks of data +/// as Strings. +/// +/// Captures stderr by default. Use `capture_stderr(False)` to override. +pub fn stream(on_data callback: fn(String) -> Nil) -> Stdio { + stream_bits(fn(bits) { callback(lossy_bits_to_string(bits)) }) +} + +/// Stream raw binary data as it arrives. The callback receives chunks of +/// data as BitArrays. +/// +/// Captures stderr by default. Use `capture_stderr(False)` to override. +pub fn stream_bits(on_data callback: fn(BitArray) -> Nil) -> Stdio { + custom(Nil, fn(_, data) { callback(data) }, fn(_, _) { Nil }) } /// Collect all output and call the callback when the process finishes. -/// The callback receives the complete output as a String and the exit status code. /// -/// This is useful when you want to process all output at once rather than streaming it. +/// The callback receives the complete output as a String and the exit status +/// code. This is useful when you want to process all output at once rather +/// than streaming it. /// -/// ## Example -/// ```gleam -/// stdio.collect(fn(output, status_code) { -/// io.println("Process exited with code: " <> int.to_string(status_code)) -/// io.println("Output: " <> output) -/// }) -/// ``` -pub fn collect(on_complete callback: fn(String, Int) -> Nil) -> Stdio { - reduce("", fn(buffer, chunk) { buffer <> chunk }, callback) +/// Captures stderr by default. Use `capture_stderr(False)` to override. +pub fn collect(on_finish callback: fn(String, Int) -> Nil) -> Stdio { + custom("", fn(buf, data) { buf <> lossy_bits_to_string(data) }, callback) +} + +/// Collect all output and call the callback when the process finishes. +/// +/// The callback receives the complete output as a BitArray and the exit status +/// code. This is useful when you want to process all output at once rather +/// than streaming it. +/// +/// Captures stderr by default. Use `capture_stderr(False)` to override. +pub fn collect_bits(on_finish callback: fn(BitArray, Int) -> Nil) -> Stdio { + custom(<<>>, fn(buffer, data) { <> }, callback) } -pub fn collect_bits(on_complete callback: fn(BitArray, Int) -> Nil) -> Stdio { - reduce_bits(<<>>, fn(buffer, chunk) { <> }, callback) +/// Add a callback when the process exits with its status code. +pub fn on_exit(stdio: Stdio, on_exit callback: fn(Int) -> Nil) -> Stdio { + let on_exit = stdio.on_exit + Stdio(..stdio, on_exit: fn(state, status) { + on_exit(state, status) + callback(status) + }) } -/// Capture stderr in addition to stdout for streaming configurations. -/// Only applies to `stream` and `lines` configurations. +// -- HELPERS ----------------------------------------------------------------- + +@external(erlang, "child_process_ffi", "bits_to_string") +@external(javascript, "../child_process_ffi.mjs", "bits_to_string") +fn lossy_bits_to_string(bits: BitArray) -> String + +// -- SELECTOR ---------------------------------------------------------------- + +@target(erlang) +/// Select messages from a running process started with `spawn_raw`. +/// +/// The process must be spawned with a mode of `capture(_)` to receive data. /// /// ## Example /// ```gleam -/// stdio.lines(fn(line) { io.print(line) }) -/// |> stdio.capture_stderr() +/// let assert Ok(proc) = +/// child_process.from_file("my_program") +/// |> child_process.spawn_raw(stdio.capture(True)) +/// +/// let selector = +/// process.new_selector() +/// |> stdio.select(proc, HandleData, HandleExit) /// ``` -pub fn capture_stderr(stdio: Stdio) -> Stdio { - case stdio { - Reduce(capture_stderr: False, ..) -> Reduce(..stdio, capture_stderr: True) - _ -> stdio - } -} +@external(erlang, "child_process_ffi", "select") +pub fn select( + selector: Selector(msg), + process: Process, + on_data: fn(BitArray) -> msg, + on_exit: fn(Int) -> msg, +) -> Selector(msg) diff --git a/src/child_process_ffi.erl b/src/child_process_ffi.erl index d611875..96c1100 100644 --- a/src/child_process_ffi.erl +++ b/src/child_process_ffi.erl @@ -1,27 +1,11 @@ -module(child_process_ffi). --export([find_executable/1, run/6, spawn/6, write/2, close/1, term/1, kill/1, - os_process_id/1]). +-export([find_executable/1, run/5, spawn/5, spawn_raw/5, write/2, close/1, term/1, kill/1, + os_process_id/1, port_owner/1, select/4]). -export([newline/0, os/0, priv_directory/0, bits_to_string/1]). --export([erlang_process_pid/1]). -include_lib("kernel/include/file.hrl"). -bits_to_string(Binary) -> - bits_to_string(Binary, <<>>). - -bits_to_string(<<>>, Acc) -> - Acc; -bits_to_string(Binary, Acc) -> - case unicode:characters_to_binary(Binary) of - Chunk when is_binary(Chunk) -> - <>; - {error, Chunk, <<_Bad, Rest/binary>>} -> - bits_to_string(Rest, <>); - {incomplete, Chunk, _Rest} -> - <> - end. - %% Find an executable in PATH find_executable(Name) -> case os:find_executable( @@ -33,40 +17,96 @@ find_executable(Name) -> {ok, unicode:characters_to_binary(Path)} end. -%% Get the native newline character for the operating system -newline() -> - case os:type() of - {win32, _} -> - <<"\r\n">>; - _ -> - <<"\n">> - end. - -%% Run a process synchronously -run(Executable, Args, Env, WorkingDirectory, Stdio, OnExit) -> - {Opts, InitialState, OnData, OnFinish} = - prepare_port_options(Args, Env, WorkingDirectory, run_stdio(Stdio), OnExit), - case open_process_port(Executable, Opts) of - {ok, Port, Monitor} -> - process_loop(Port, Monitor, OnData, OnFinish, InitialState, -1); +run(Executable, Args, Env, WorkingDirectory, Mode) -> + OnData = fun(Acc, Data) -> <> end, + OnFinish = fun(Data, Code) -> {ok, {output, Code, Data}} end, + case spawn_raw(Executable, Args, Env, WorkingDirectory, Mode) of + {ok, Port} -> + Monitor = erlang:monitor(port, Port), + loop(Port, Monitor, undefined, OnData, OnFinish, <<>>, -1, normal); {error, Error} -> {error, Error} end. -run_stdio(inherit) -> - inherit; -run_stdio(_) -> - OnData = fun(Chunks, Data) -> [Chunks, Data] end, - OnFinish = fun(Chunks, Code) -> {ok, {output, Code, iolist_to_binary(Chunks)}} end, - {reduce, [], OnData, OnFinish, true}. +spawn(Executable, Args, Env, WorkingDirectory, Stdio) -> + {stdio, Mode, InitialState, OnData, OnFinish} = Stdio, + Self = self(), + Ref = make_ref(), + + Proc = + fun() -> + process_flag(trap_exit, true), + case spawn_raw(Executable, Args, Env, WorkingDirectory, Mode) of + {ok, Port} -> + Monitor = erlang:monitor(port, Port), + Self ! {Ref, ok, Port}, + loop(Port, Monitor, Self, OnData, OnFinish, InitialState, -1, normal); + {error, Error} -> Self ! {Ref, error, Error} + end + end, + + Pid = proc_lib:spawn_link(Proc), + Monitor = erlang:monitor(process, Pid), + + receive + {Ref, ok, Port} -> + erlang:demonitor(Monitor, [flush]), + {ok, Port}; + {Ref, error, Error} -> + erlang:demonitor(Monitor, [flush]), + {error, Error}; + {'DOWN', Monitor, process, Pid, _Reason} -> + {error, runtime_limit_reached} + end. + +%% Get the Pid of the process that currently owns the port. +port_owner(Port) -> + case catch erlang:port_info(Port, connected) of + {connected, Pid} -> + {ok, Pid}; + _ -> + {error, nil} + end. -%% Prepare port options and callbacks from process configuration -prepare_port_options(Args, Env, WorkingDirectory, Stdio, OnExit) -> +select(Selector, Port, OnData, OnExit) -> + Handler = + fun({_Port, Msg}) -> + case Msg of + {data, Chunk} when is_binary(Chunk) -> OnData(Chunk); + {exit_status, ExitCode} when is_integer(ExitCode) -> OnExit(ExitCode) + end + end, + + gleam@erlang@process:select_record(Selector, Port, 2, Handler). + +spawn_raw(Executable, Args, Env, WorkingDirectory, Mode) -> + Opts = prepare_port_options(Args, Env, WorkingDirectory, Mode), + PortName = {spawn_executable, unicode:characters_to_list(Executable)}, + try + {ok, erlang:open_port(PortName, Opts)} + catch + error:system_limit -> + {error, runtime_limit_reached}; + error:enomem -> + {error, out_of_memory}; + error:eagain -> + {error, os_process_limit_reached}; + error:enfile -> + {error, os_file_limit_reached}; + error:enametoolong -> + {error, command_too_long}; + error:emfile -> + {error, not_enough_file_descriptors}; + error:eacces -> + {error, {file_not_executable, Executable}}; + error:enoent -> + {error, {file_not_found, Executable}} + end. + +prepare_port_options(Args, Env, WorkingDirectory, Stdio) -> BaseOptions = [{args, Args}, {env, merge_env(Env)}, exit_status, binary, hide], - Options = stdio_open_options(cwd_options(BaseOptions, WorkingDirectory), Stdio), - std_callbacks(Options, Stdio, OnExit). + cwd_options(stdio_open_options(BaseOptions, Stdio), WorkingDirectory). -%% Merge environment variables with current process environment merge_env([]) -> []; merge_env(Env) -> @@ -87,96 +127,28 @@ stdio_open_options(Opts, null) -> [use_stdio, stderr_to_stdout, stream | Opts]; stdio_open_options(Opts, inherit) -> [nouse_stdio | Opts]; -stdio_open_options(Opts, {reduce, _, _, _, true}) -> +stdio_open_options(Opts, {capture, true}) -> [use_stdio, stderr_to_stdout, stream | Opts]; -stdio_open_options(Opts, _) -> +stdio_open_options(Opts, {capture, false}) -> [use_stdio, stream | Opts]. -std_callbacks(Opts, {reduce, Initial, OnDataFn, OnFinishFn, _}, OnExit) -> - WrappedExit = - fun(Acc, Code) -> - Result = OnFinishFn(Acc, Code), - OnExit(Code), - Result - end, - {Opts, Initial, OnDataFn, WrappedExit}; -std_callbacks(Opts, _Stdio, OnExit) -> - WrappedExit = - fun(_, Code) -> - OnExit(Code), - {ok, {output, Code, <<>>}} - end, - {Opts, nil, fun(_, _) -> nil end, WrappedExit}. - -%% Spawn a process asynchronously -spawn(Executable, Args, Env, WorkingDirectory, Stdio, OnExit) -> - {Opts, InitialAcc, OnData, OnFinish} = - prepare_port_options(Args, Env, WorkingDirectory, Stdio, OnExit), - - Self = self(), - Ref = make_ref(), - - ProcFn = - fun() -> process_init(Self, Ref, Executable, Opts, InitialAcc, OnData, OnFinish) end, - - Pid = proc_lib:spawn_link(ProcFn), - % Wait for port to be opened - receive - {Ref, port_opened, Port} -> - {ok, {Port, Pid}}; - {Ref, port_error, Error} -> - {error, Error} - end. - -process_init(Parent, Ref, Executable, PortSettings, InitialState, OnData, OnFinish) -> - case open_process_port(Executable, PortSettings) of - {ok, Port, Monitor} -> - Parent ! {Ref, port_opened, Port}, - process_loop(Port, Monitor, OnData, OnFinish, InitialState, -1); - {error, Error} -> - Parent ! {Ref, port_error, Error} - end. - -%% Open port and return it with monitor - caller decides what to do next -open_process_port(Executable, Opts) -> - PortName = {spawn_executable, unicode:characters_to_list(Executable)}, - case catch erlang:open_port(PortName, Opts) of - {'EXIT', {enoent, _}} -> - {error, {file_not_found, Executable}}; - {'EXIT', {eacces, _}} -> - {error, {file_not_executable, Executable}}; - {'EXIT', _} -> - {error, system_limit}; - Port -> - {ok, Port, erlang:monitor(port, Port)} - end. - -process_loop(Port, Monitor, OnData, OnFinish, State, StatusCode) -> +loop(Port, Monitor, Parent, OnData, OnFinish, State, StatusCode, ExitReason) -> receive {Port, {data, Chunk}} -> - %% Call on_data callback NewState = OnData(State, Chunk), - process_loop(Port, Monitor, OnData, OnFinish, NewState, StatusCode); + loop(Port, Monitor, Parent, OnData, OnFinish, NewState, StatusCode, ExitReason); {Port, {exit_status, Code}} -> - %% Process exited, just store the status code and wait for DOWN - process_loop(Port, Monitor, OnData, OnFinish, State, Code); - {'DOWN', Monitor, port, Port, normal} when StatusCode =:= -1 -> - flush_exit(Port), - OnFinish(State, 0); + loop(Port, Monitor, Parent, OnData, OnFinish, State, Code, ExitReason); {'DOWN', Monitor, port, Port, _} -> flush_exit(Port), - OnFinish(State, StatusCode); - close -> - %% Close the port. On Erlang, this closes stdin, stdout, and stderr. - %% This is a platform limitation - we cannot close just stdin. - %% Any buffered output from the child that hasn't been read yet will be lost. - erlang:port_close(Port), - process_loop(Port, Monitor, OnData, OnFinish, State, StatusCode) + loop_exit(OnFinish(State, StatusCode), ExitReason); + {'EXIT', Port, _Reason} -> + loop(Port, Monitor, Parent, OnData, OnFinish, State, StatusCode, ExitReason); + {'EXIT', Parent, Reason} -> + term(Port), + loop(Port, Monitor, Parent, OnData, OnFinish, State, StatusCode, Reason) end. -% If the calling process traps exits, this will remove that exit message -% from the inbox. Since we call it after receiving DOWN, we know for sure that -% the exit message would also be in the inbox by now. flush_exit(Port) -> receive {'EXIT', Port, _} -> @@ -185,47 +157,91 @@ flush_exit(Port) -> ok end. +loop_exit(State, normal) -> + State; +loop_exit(_State, {shutdown, Reason}) -> + exit(Reason); +loop_exit(_State, Reason) -> + exit(Reason). + %% Write text to process stdin -write({Port, _Pid}, Text) -> - erlang:port_command(Port, Text), - nil. +write(Port, Text) -> + try erlang:port_command(Port, Text) of + true -> + {ok, nil}; + false -> + {error, write_aborted} + catch + error:badarg -> + {error, process_exited} + end. -%% Close process stdin -close({_Port, Pid}) -> - Pid ! close, +%% Close process stdin. +close(Port) -> + catch erlang:port_close(Port), nil. %% Stop process gracefully (SIGTERM) -term({Port, _Pid}) -> - os:cmd(stop_command(os:type(), Port)), - nil. +term(Port) -> + case os_process_id(Port) of + {ok, OsPid} -> + os:cmd(stop_command(os:type(), OsPid)), + nil; + {error, nil} -> + nil + end. -stop_command({win32, _}, Port) -> - io_lib:format("taskkill /PID ~p", [os_pid(Port)]); -stop_command(_, Port) -> - io_lib:format("kill -TERM ~p", [os_pid(Port)]). +stop_command({win32, _}, OsPid) -> + io_lib:format("taskkill /PID ~p", [OsPid]); +stop_command(_, OsPid) -> + io_lib:format("kill -TERM ~p", [OsPid]). %% Kill process forcefully (SIGKILL) -kill({Port, _Pid}) -> - os:cmd(kill_command(os:type(), Port)), - nil. +kill(Port) -> + case os_process_id(Port) of + {ok, OsPid} -> + os:cmd(kill_command(os:type(), OsPid)), + nil; + {error, nil} -> + nil + end. -kill_command({win32, _}, Port) -> - io_lib:format("taskkill /F /PID ~p", [os_pid(Port)]); -kill_command(_, Port) -> - io_lib:format("kill -KILL ~p", [os_pid(Port)]). +kill_command({win32, _}, OsPid) -> + io_lib:format("taskkill /F /PID ~p", [OsPid]); +kill_command(_, OsPid) -> + io_lib:format("kill -KILL ~p", [OsPid]). -%% Get OS process ID -os_process_id({Port, _Pid}) -> - os_pid(Port). +os_process_id(Port) -> + case catch erlang:port_info(Port, os_pid) of + {os_pid, OsPid} -> + {ok, OsPid}; + _ -> + {error, nil} + end. -erlang_process_pid({_Port, Pid}) -> - Pid. +bits_to_string(Binary) -> + bits_to_string(Binary, <<>>). -%% Internal helper to get OS PID from port -os_pid(Port) -> - {os_pid, OsPid} = erlang:port_info(Port, os_pid), - OsPid. +bits_to_string(<<>>, Acc) -> + Acc; +bits_to_string(Binary, Acc) -> + case unicode:characters_to_binary(Binary) of + Chunk when is_binary(Chunk) -> + <>; + {error, Chunk, <<_Bad, Rest/binary>>} -> + bits_to_string(Rest, <>); + {incomplete, Chunk, _Rest} -> + <> + end. + +%% Get the native newline character for the operating system +newline() -> + case os:type() of + {win32, _} -> + <<"\r\n">>; + _ -> + <<"\n">> + end. %% Detect the operating system os() -> diff --git a/src/child_process_ffi.mjs b/src/child_process_ffi.mjs index e66e629..7372b0f 100644 --- a/src/child_process_ffi.mjs +++ b/src/child_process_ffi.mjs @@ -4,24 +4,31 @@ import fs from "node:fs"; import pathModule from "node:path"; import { fileURLToPath } from "node:url"; -import { Result$Ok, Result$Error, BitArray$BitArray } from "./gleam.mjs"; +import { + Result$Ok, + Result$Error, + Result$isOk, + Result$Ok$0, + BitArray$BitArray, +} from "./gleam.mjs"; import * as List from "../gleam_stdlib/gleam/list.mjs"; import { Option$isSome, Option$Some$0 } from "../gleam_stdlib/gleam/option.mjs"; import { StartError$FileNotFound, StartError$FileNotExecutable, - StartError$SystemLimit, + StartError$RuntimeLimitReached, Output$Output, } from "../child_process/child_process.mjs"; import { - Stdio$isInherit, - Stdio$isNull, - Stdio$isReduce, - Stdio$Reduce$on_data, - Stdio$Reduce$on_exit, - Stdio$Reduce$capture_stderr, - Stdio$Reduce$initial, + Mode$isInherit, + Mode$isNull, + Mode$isCapture, + Mode$Capture$capture_stderr, + Stdio$Stdio$mode, + Stdio$Stdio$initial, + Stdio$Stdio$on_data, + Stdio$Stdio$on_exit, Os$Windows, Os$Macos, Os$Linux, @@ -165,43 +172,71 @@ export function newline() { return isWindows ? "\r\n" : "\n"; } -// Run a process synchronously -export function run(executable, args, envVars, cwd, stdio, onExit) { +export function run(executable, args, env, cwd, mode) { const options = { - env: getEnv(envVars), - - cwd: (Option$isSome(cwd) && Option$Some$0(cwd)) || undefined, - - encoding: Stdio$isInherit(stdio) ? "buffer" : "utf8", - stdio: Stdio$isInherit(stdio) ? "inherit" : ["ignore", "pipe", "pipe"], - + env: getEnv(env), + cwd: Option$isSome(cwd) ? Option$Some$0(cwd) : undefined, + encoding: Mode$isInherit(mode) ? "buffer" : "utf8", + stdio: Mode$isInherit(mode) ? "inherit" : ["ignore", "pipe", "pipe"], windowsHide: true, maxBuffer: Infinity, }; const result = child_process.spawnSync(executable, toArray(args), options); if (result.error) { - return handleError(executable, result.error); + switch (result.error.code) { + case "ENOENT": + return Result$Error(StartError$FileNotFound(executable)); + case "EACESS": + return Result$Error(StartError$FileNotExecutable(executable)); + case "ENOMEM": + return Result$Error(StartError$OutOfMemory()); + case "EAGAIN": + return Result$Error(StartError$OsProcessLimitReached()); + case "ENFILE": + return Result$Error(StartError$OsFileLimitReached()); + case "ENAMETOOLONG": + return Result$Error(StartError$CommandTooLong()); + case "EMFILE": + return Result$Error(StartError$NotEnoughFileDescriptors()); + default: + return Result$Error(StartError$RuntimeLimitReached()); + } } // I have not found a good way to mix stdout and stderr in this case. const output = (result.stdout || "") + (result.stderr || ""); const status_code = result.status ?? -1; - onExit(status_code); - return Result$Ok(Output$Output(status_code, output)); } -// Spawn a process asynchronously -export function spawn( - executable, - args, - envVars, - workingDirectory, - stdioConfig, - onExit, -) { +export function spawn(executable, args, env, cwd, stdio) { + const result = spawn_raw(executable, args, env, cwd, Stdio$Stdio$mode(stdio)); + if (!Result$isOk(result)) { + return result; + } + + const process = Result$Ok$0(result); + let state = Stdio$Stdio$initial(stdio); + + function handleData(chunk) { + const bits = BitArray$BitArray(new Uint8Array(chunk)); + state = Stdio$Stdio$on_data(stdio)(state, bits); + } + function handleExit(code) { + Stdio$Stdio$on_exit(stdio)(state, code ?? -1); + } + + process.stdout?.on("data", handleData); + process.stderr?.on("data", handleData); + process.on("exit", handleExit); + process.on("error", () => handleExit(-2)); + + return result; +} + +export function spawn_raw(executable, args, envVars, workingDirectory, mode) { const stat = fs.statSync(executable, { throwIfNoEntry: false }); if (stat == null || !stat.isFile()) { return Result$Error(StartError$FileNotFound(executable)); @@ -210,9 +245,9 @@ export function spawn( return Result$Error(StartError$FileNotExecutable(executable)); } - const cwd = - (Option$isSome(workingDirectory) && Option$Some$0(workingDirectory)) || - undefined; + const cwd = Option$isSome(workingDirectory) + ? Option$Some$0(workingDirectory) + : undefined; if (cwd !== undefined) { const cwdStat = fs.statSync(cwd, { throwIfNoEntry: false }); if (!cwdStat || !cwdStat.isDirectory) { @@ -224,50 +259,25 @@ export function spawn( const env = getEnv(envVars); let stdio; - if (Stdio$isInherit(stdioConfig)) { + if (Mode$isInherit(mode)) { stdio = "inherit"; - } else if (Stdio$isNull(stdioConfig)) { + } else if (Mode$isNull(mode)) { stdio = ["pipe", "ignore", "ignore"]; // stdin is pipe to allow writing - } else if (Stdio$isReduce(stdioConfig)) { - const stderr = Stdio$Reduce$capture_stderr(stdioConfig) ? "pipe" : "inherit"; + } else if (Mode$isCapture(mode)) { + const stderr = Mode$Capture$capture_stderr(mode) ? "pipe" : "inherit"; stdio = ["pipe", "pipe", stderr]; } else { - throw new Error("invalid stdio"); + throw new Error("invalid mode"); } const options = { env, cwd, stdio, windowsHide: true }; - const child = child_process.spawn(executable, toArray(args), options); - - let acc; - if (Stdio$isReduce(stdioConfig)) { - acc = Stdio$Reduce$initial(stdioConfig); - - function handleStdio(chunk) { - acc = Stdio$Reduce$on_data(stdioConfig)(acc, chunk.toString("utf8")); - } - - child.stdout?.on("data", handleStdio); - child.stderr?.on("data", handleStdio); - } - - function handleExit(code) { - if (!onExit) return; - code ??= -1; - - // Call on_exit callback with final accumulator state - if (Stdio$isReduce(stdioConfig)) { - Stdio$Reduce$on_exit(stdioConfig)(acc, code); - } - - onExit?.(code); - onExit = null; + try { + const child = child_process.spawn(executable, toArray(args), options); + return Result$Ok(child); + } catch { + return Result$Error(StartError$RuntimeLimitReached()); } - - child.on("error", () => handleExit(-2)); - child.on("exit", handleExit); - - return Result$Ok(child); } function getEnv(env) { @@ -281,39 +291,53 @@ function getEnv(env) { }); } -function handleError(executable, error) { - if (error.code === "ENOENT") { - return Result$Error(StartError$FileNotFound(executable)); - } else if (error.code === "EACCES") { - return Result$Error(StartError$FileNotExecutable(executable)); - } else { - return Result$Error(StartError$SystemLimit()); - } -} - // Write text to process stdin export function write(process, text) { - process.stdin.write(text); + if (process.exitCode !== null || process.stdin.destroyed) { + return Result$Error(WriteError$ProcessExited()); + } + + if (process.stdin.write(text)) { + return Result$Ok(); + } else { + return Result$Error(WriteError$WriteAborted()); + } } // Close process stdin export function close(process) { + if (process.exitCode !== null || process.stdin.destroyed) { + return; + } + process.stdin.end(); } // Stop process gracefully (SIGTERM) export function term(process) { + if (process.exitCode !== null) { + return; + } + process.kill("SIGTERM"); } // Kill process forcefully (SIGKILL) export function kill(process) { + if (process.exitCode !== null) { + return; + } + process.kill("SIGKILL"); } // Get OS process ID export function os_process_id(process) { - return process.pid; + if (process.exitCode !== null) { + return Result$Error(); + } + + return Result$Ok(process.pid); } // Detect the operating system diff --git a/test/child_process_test.gleam b/test/child_process_test.gleam index 00e0ea1..10fb504 100644 --- a/test/child_process_test.gleam +++ b/test/child_process_test.gleam @@ -1,4 +1,5 @@ import child_process +import child_process/stdio import gleam/list import gleam/string import gleeunit @@ -10,9 +11,9 @@ pub fn main() -> Nil { // Basic run tests pub fn simple_echo_test() { let assert Ok(output) = - child_process.new("/bin/echo") + child_process.from_file("/bin/echo") |> child_process.arg("hello world") - |> child_process.run() + |> child_process.run(stdio.capture(True)) assert output.output == "hello world\n" assert output.status_code == 0 @@ -20,18 +21,18 @@ pub fn simple_echo_test() { pub fn exit_code_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "exit 42"]) - |> child_process.run() + |> child_process.run(stdio.null()) assert output.status_code == 42 } pub fn large_output_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "for i in $(seq 1 100); do echo line_$i; done"]) - |> child_process.run() + |> child_process.run(stdio.capture(True)) let line_count = output.output @@ -44,9 +45,9 @@ pub fn large_output_test() { pub fn line_without_newline_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "printf 'no newline'"]) - |> child_process.run() + |> child_process.run(stdio.capture(True)) assert output.output == "no newline" } @@ -54,16 +55,16 @@ pub fn line_without_newline_test() { // Error handling pub fn file_not_found_test() { let assert Error(child_process.FileNotFound(_)) = - child_process.new("/nonexistent/executable") - |> child_process.run() + child_process.from_file("/nonexistent/executable") + |> child_process.run(stdio.null()) } // Working directory pub fn working_directory_test() { let assert Ok(output) = - child_process.new("/bin/pwd") + child_process.from_file("/bin/pwd") |> child_process.cwd("/tmp") - |> child_process.run() + |> child_process.run(stdio.capture(True)) // On macOS, /tmp is a symlink to /private/tmp assert output.output == "/tmp\n" || output.output == "/private/tmp\n" @@ -82,67 +83,67 @@ pub fn shell_test() { pub fn environment_variable_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "echo $MY_VAR"]) |> child_process.env("MY_VAR", "test_value") - |> child_process.run() + |> child_process.run(stdio.capture(True)) assert output.output == "test_value\n" } pub fn spawn_basic_test() { let assert Ok(proc) = - child_process.new("/bin/sleep") + child_process.from_file("/bin/sleep") |> child_process.arg("0.1") - |> child_process.spawn() + |> child_process.spawn(stdio.ignore()) - let pid = child_process.os_process_id(proc) + let assert Ok(pid) = child_process.os_process_id(proc) assert pid != 0 } pub fn stop_process_test() { let assert Ok(proc) = - child_process.new("/bin/sleep") + child_process.from_file("/bin/sleep") |> child_process.arg("10") - |> child_process.spawn() + |> child_process.spawn(stdio.ignore()) child_process.stop(proc) } pub fn kill_process_test() { let assert Ok(proc) = - child_process.new("/bin/sleep") + child_process.from_file("/bin/sleep") |> child_process.arg("10") - |> child_process.spawn() + |> child_process.spawn(stdio.ignore()) child_process.kill(proc) } pub fn spawn_with_write_test() { let assert Ok(proc) = - child_process.new("/bin/cat") - |> child_process.spawn() + child_process.from_file("/bin/cat") + |> child_process.spawn(stdio.ignore()) - child_process.write(proc, "hello\n") - child_process.write(proc, "world\n") + let _ = child_process.write(proc, "hello\n") + let _ = child_process.write(proc, "world\n") child_process.close(proc) } pub fn multiple_environment_variables_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "echo $VAR1-$VAR2"]) |> child_process.envs([#("VAR1", "hello"), #("VAR2", "world")]) - |> child_process.run() + |> child_process.run(stdio.capture(True)) assert output.output == "hello-world\n" } pub fn lines_mode_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "echo 'line1' && echo 'line2'"]) - |> child_process.run() + |> child_process.run(stdio.capture(True)) assert string.contains(output.output, "line1") } @@ -160,9 +161,9 @@ pub fn find_executable_not_found_test() { pub fn empty_output_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", ""]) - |> child_process.run() + |> child_process.run(stdio.null()) assert output.output == "" assert output.status_code == 0 @@ -170,9 +171,9 @@ pub fn empty_output_test() { pub fn process_with_output_then_exit_test() { let assert Ok(output) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "echo 'done' && exit 5"]) - |> child_process.run() + |> child_process.run(stdio.capture(True)) assert output.status_code == 5 assert output.output == "done\n" diff --git a/test/on_exit_test.gleam b/test/on_exit_test.gleam index 0bf9fbe..0206ca9 100644 --- a/test/on_exit_test.gleam +++ b/test/on_exit_test.gleam @@ -1,8 +1,9 @@ +import gleeunit + +@target(erlang) import child_process @target(erlang) import child_process/stdio -import gleeunit - @target(erlang) import gleam/erlang/process @@ -10,46 +11,17 @@ pub fn main() -> Nil { gleeunit.main() } -// Test OnExit callback is called with correct status code for run() -@target(erlang) -pub fn run_on_exit_success_test() { - let subject = process.new_subject() - - let assert Ok(output) = - child_process.new("/bin/sh") - |> child_process.args(["-c", "exit 0"]) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) - |> child_process.run() - - assert output.status_code == 0 - assert process.receive(subject, 100) == Ok(0) -} - -@target(erlang) -pub fn run_on_exit_failure_test() { - let subject = process.new_subject() - - let assert Ok(output) = - child_process.new("/bin/sh") - |> child_process.args(["-c", "exit 7"]) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) - |> child_process.run() - - assert output.status_code == 7 - assert process.receive(subject, 100) == Ok(7) -} - -// Test OnExit callback is called with correct status code for spawn() +// Test exit handler is called with correct status code after spawn @target(erlang) pub fn spawn_on_exit_success_test() { let subject = process.new_subject() - let assert Ok(_process) = - child_process.new("/bin/sh") + let assert Ok(_proc) = + child_process.from_file("/bin/sh") |> child_process.args(["-c", "exit 0"]) - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject, code) }), + ) assert process.receive(subject, 1000) == Ok(0) } @@ -58,117 +30,100 @@ pub fn spawn_on_exit_success_test() { pub fn spawn_on_exit_failure_test() { let subject = process.new_subject() - let assert Ok(_process) = - child_process.new("/bin/sh") + let assert Ok(_proc) = + child_process.from_file("/bin/sh") |> child_process.args(["-c", "exit 13"]) - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject, code) }), + ) assert process.receive(subject, 1000) == Ok(13) } -// Test OnExit with stdio.collect() +// Test collector handler after spawn @target(erlang) -pub fn collect_on_exit_test() { - let exit_subject = process.new_subject() - let output_subject = process.new_subject() +pub fn collect_test() { + let subject = process.new_subject() - let assert Ok(_process) = - child_process.new("/bin/echo") + let assert Ok(_proc) = + child_process.from_file("/bin/echo") |> child_process.arg("hello") - |> child_process.stdio( + |> child_process.spawn( stdio.collect(fn(output, status) { - process.send(output_subject, #(output, status)) + process.send(subject, #(output, status)) }), ) - |> child_process.on_exit(fn(code) { process.send(exit_subject, code) }) - |> child_process.spawn() - // Both callbacks should be called - assert process.receive(exit_subject, 1000) == Ok(0) - assert process.receive(output_subject, 1000) == Ok(#("hello\n", 0)) + assert process.receive(subject, 1000) == Ok(#("hello\n", 0)) } -// Test OnExit with stdio.stream() +// Test chunks handler after spawn @target(erlang) -pub fn stream_on_exit_test() { - let exit_subject = process.new_subject() - let data_subject = process.new_subject() +pub fn stream_test() { + let subject = process.new_subject() - let assert Ok(_process) = - child_process.new("/bin/echo") + let assert Ok(_proc) = + child_process.from_file("/bin/echo") |> child_process.arg("test") - |> child_process.stdio( - stdio.stream(fn(chunk) { process.send(data_subject, chunk) }), + |> child_process.spawn( + stdio.stream(fn(chunk) { process.send(subject, chunk) }), ) - |> child_process.on_exit(fn(code) { process.send(exit_subject, code) }) - |> child_process.spawn() - // Should receive data first, then exit - assert process.receive(data_subject, 1000) == Ok("test\n") - assert process.receive(exit_subject, 1000) == Ok(0) + assert process.receive(subject, 1000) == Ok("test\n") } -// Test OnExit with stdio.lines() +// Test lines handler after spawn @target(erlang) -pub fn lines_on_exit_test() { - let exit_subject = process.new_subject() - let lines_subject = process.new_subject() +pub fn lines_test() { + let subject = process.new_subject() - let assert Ok(_process) = - child_process.new("/bin/sh") + let assert Ok(_proc) = + child_process.from_file("/bin/sh") |> child_process.args(["-c", "echo line1; echo line2"]) - |> child_process.stdio( - stdio.lines(fn(line) { process.send(lines_subject, line) }), + |> child_process.spawn( + stdio.lines(fn(line) { process.send(subject, line) }), ) - |> child_process.on_exit(fn(code) { process.send(exit_subject, code) }) - |> child_process.spawn() - // Should receive both lines, then exit - assert process.receive(lines_subject, 1000) == Ok("line1\n") - assert process.receive(lines_subject, 1000) == Ok("line2\n") - assert process.receive(exit_subject, 1000) == Ok(0) + assert process.receive(subject, 1000) == Ok("line1\n") + assert process.receive(subject, 1000) == Ok("line2\n") } -// Test OnExit is called even when process is stopped +// Test exit handler is called even when process is stopped @target(erlang) pub fn on_exit_after_stop_test() { let subject = process.new_subject() let assert Ok(proc) = - child_process.new("/bin/sleep") + child_process.from_file("/bin/sleep") |> child_process.arg("10") - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject, code) }), + ) child_process.stop(proc) // Should receive SIGTERM exit code (143) let assert Ok(code) = process.receive(subject, 1000) assert code == 143 || code == 15 - // Different platforms may report differently } -// Test OnExit is called even when process is killed +// Test exit handler is called even when process is killed @target(erlang) pub fn on_exit_after_kill_test() { let subject = process.new_subject() let assert Ok(proc) = - child_process.new("/bin/sleep") + child_process.from_file("/bin/sleep") |> child_process.arg("10") - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject, code) }), + ) child_process.kill(proc) // Should receive SIGKILL exit code (137) let assert Ok(code) = process.receive(subject, 1000) assert code == 137 || code == 9 - // Different platforms may report differently } // Test multiple processes with different exit codes @@ -179,55 +134,42 @@ pub fn multiple_on_exit_test() { let subject3 = process.new_subject() let assert Ok(_p1) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "exit 1"]) - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject1, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject1, code) }), + ) let assert Ok(_p2) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "exit 2"]) - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject2, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject2, code) }), + ) let assert Ok(_p3) = - child_process.new("/bin/sh") + child_process.from_file("/bin/sh") |> child_process.args(["-c", "exit 3"]) - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject3, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject3, code) }), + ) - // Each should receive its own exit code assert process.receive(subject1, 1000) == Ok(1) assert process.receive(subject2, 1000) == Ok(2) assert process.receive(subject3, 1000) == Ok(3) } -// Test OnExit with null stdio +// Test exit handler with null stdio (ignore) @target(erlang) pub fn null_stdio_on_exit_test() { let subject = process.new_subject() - let assert Ok(_process) = - child_process.new("/bin/echo") + let assert Ok(_proc) = + child_process.from_file("/bin/echo") |> child_process.arg("discarded") - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) - |> child_process.spawn() + |> child_process.spawn( + stdio.ignore() |> stdio.on_exit(fn(code) { process.send(subject, code) }), + ) assert process.receive(subject, 1000) == Ok(0) } - -// Test that OnExit callback with no-op still works -pub fn noop_on_exit_test() { - let assert Ok(output) = - child_process.new("/bin/echo") - |> child_process.arg("test") - |> child_process.on_exit(fn(_code) { Nil }) - |> child_process.run() - - assert output.output == "test\n" - assert output.status_code == 0 -} diff --git a/test/stream_test.gleam b/test/stream_test.gleam new file mode 100644 index 0000000..84b2768 --- /dev/null +++ b/test/stream_test.gleam @@ -0,0 +1,17 @@ +import child_process +import child_process/stdio +import gleam/erlang/process + +pub fn main() { + let assert Ok(_) = + child_process.from_file("test/stream_test.sh") + |> child_process.spawn( + stdio.stream(fn(chunk) { + echo chunk + Nil + }), + ) + + process.sleep(5000) + echo "main exit" +} diff --git a/test/stream_test.sh b/test/stream_test.sh new file mode 100755 index 0000000..ee4ef88 --- /dev/null +++ b/test/stream_test.sh @@ -0,0 +1,5 @@ +#!/bin/sh +while true; do + echo $(date) + sleep 1 +done diff --git a/test/supervision_test.gleam b/test/supervision_test.gleam index c278d9a..983b24d 100644 --- a/test/supervision_test.gleam +++ b/test/supervision_test.gleam @@ -1,10 +1,17 @@ +import gleeunit + +@target(erlang) import child_process +@target(erlang) import child_process/stdio +@target(erlang) import gleam/erlang/process +@target(erlang) import gleam/otp/actor +@target(erlang) import gleam/otp/factory_supervisor +@target(erlang) import gleam/otp/static_supervisor as supervisor -import gleeunit pub fn main() -> Nil { gleeunit.main() @@ -14,9 +21,8 @@ pub fn main() -> Nil { @target(erlang) pub fn supervised_creates_spec_test() { let _spec = - child_process.new("/bin/cat") - |> child_process.stdio(stdio.null()) - |> child_process.supervised() + child_process.from_file("/bin/cat") + |> child_process.supervised(stdio: stdio.ignore()) Nil } @@ -25,10 +31,9 @@ pub fn supervised_creates_spec_test() { @target(erlang) pub fn supervised_with_supervisor_test() { let child_spec = - child_process.new("/bin/sleep") + child_process.from_file("/bin/sleep") |> child_process.arg("10") - |> child_process.stdio(stdio.null()) - |> child_process.supervised() + |> child_process.supervised(stdio: stdio.ignore()) let assert Ok(actor.Started(pid, _)) = supervisor.new(supervisor.OneForOne) @@ -38,7 +43,6 @@ pub fn supervised_with_supervisor_test() { process.sleep(100) assert process.is_alive(pid) == True - // Proper OTP shutdown using send_exit process.send_exit(pid) process.sleep(100) assert process.is_alive(pid) == False @@ -47,7 +51,7 @@ pub fn supervised_with_supervisor_test() { // Test factory() creates a factory supervisor builder @target(erlang) pub fn factory_creates_builder_test() { - let _factory_builder = child_process.factory() + let _factory_builder = child_process.factory(stdio: stdio.ignore()) Nil } @@ -55,7 +59,7 @@ pub fn factory_creates_builder_test() { @target(erlang) pub fn factory_basic_test() { let assert Ok(actor.Started(factory_pid, _factory)) = - child_process.factory() + child_process.factory(stdio: stdio.ignore()) |> factory_supervisor.start() assert process.is_alive(factory_pid) == True @@ -68,58 +72,42 @@ pub fn factory_basic_test() { // Test factory() can spawn children dynamically @target(erlang) pub fn factory_spawn_child_test() { - let subject = process.new_subject() - let assert Ok(actor.Started(factory_pid, factory)) = - child_process.factory() + child_process.factory(stdio: stdio.ignore()) |> factory_supervisor.start() let child_config = - child_process.new("/bin/echo") + child_process.from_file("/bin/echo") |> child_process.arg("hello") - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject, code) }) let assert Ok(actor.Started(_child_pid, _child)) = factory_supervisor.start_child(factory, child_config) - // Child should exit successfully - assert process.receive(subject, 1000) == Ok(0) - + process.sleep(200) process.send_exit(factory_pid) } // Test factory() can spawn multiple children @target(erlang) pub fn factory_spawn_multiple_test() { - let subject1 = process.new_subject() - let subject2 = process.new_subject() - let assert Ok(actor.Started(factory_pid, factory)) = - child_process.factory() + child_process.factory(stdio: stdio.ignore()) |> factory_supervisor.start() let child1_config = - child_process.new("/bin/echo") + child_process.from_file("/bin/echo") |> child_process.arg("child1") - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject1, code) }) let child2_config = - child_process.new("/bin/echo") + child_process.from_file("/bin/echo") |> child_process.arg("child2") - |> child_process.stdio(stdio.null()) - |> child_process.on_exit(fn(code) { process.send(subject2, code) }) let assert Ok(actor.Started(_, _)) = factory_supervisor.start_child(factory, child1_config) let assert Ok(actor.Started(_, _)) = factory_supervisor.start_child(factory, child2_config) - // Both children should exit successfully - assert process.receive(subject1, 1000) == Ok(0) - assert process.receive(subject2, 1000) == Ok(0) - + process.sleep(200) process.send_exit(factory_pid) } @@ -127,14 +115,11 @@ pub fn factory_spawn_multiple_test() { @target(erlang) pub fn factory_handles_failure_test() { let assert Ok(actor.Started(factory_pid, factory)) = - child_process.factory() + child_process.factory(stdio: stdio.ignore()) |> factory_supervisor.start() - let bad_config = - child_process.new("/nonexistent/executable") - |> child_process.stdio(stdio.null()) + let bad_config = child_process.from_file("/nonexistent/executable") - // Should fail to start let result = factory_supervisor.start_child(factory, bad_config) case result { diff --git a/test/vm_exit_test.gleam b/test/vm_exit_test.gleam new file mode 100644 index 0000000..b1cb498 --- /dev/null +++ b/test/vm_exit_test.gleam @@ -0,0 +1,48 @@ +@target(erlang) +import child_process +@target(erlang) +import child_process/stdio +@target(erlang) +import gleam/erlang/process +@target(erlang) +import gleam/int +@target(erlang) +import gleam/io + +@target(erlang) +pub fn main() { + io.println("Spawning long-running processes...") + + let assert Ok(proc1) = + child_process.from_file("/bin/sleep") + |> child_process.arg("300") + |> child_process.spawn(stdio.ignore()) + + let assert Ok(proc2) = + child_process.from_file("/bin/sleep") + |> child_process.arg("300") + |> child_process.spawn(stdio.ignore()) + + let assert Ok(proc3) = + child_process.from_file("/bin/sleep") + |> child_process.arg("300") + |> child_process.spawn(stdio.ignore()) + + let assert Ok(pid1) = child_process.os_process_id(proc1) + let assert Ok(pid2) = child_process.os_process_id(proc2) + let assert Ok(pid3) = child_process.os_process_id(proc3) + + io.println("PID:" <> int.to_string(pid1)) + io.println("PID:" <> int.to_string(pid2)) + io.println("PID:" <> int.to_string(pid3)) + + // Wait for processes to start + process.sleep(500) + + io.println("READY") + + // Wait a bit, then exit + process.sleep(1000) + + io.println("EXITING") +}