Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,12 @@ def pid(self) -> int:

@property
def stdout(self) -> IO[bytes] | None:
"""Combined binary output pipe."""
"""Binary standard-output pipe when requested."""
...

@property
def stderr(self) -> IO[bytes] | None:
"""Binary standard-error pipe when requested separately."""
...

def kill(self) -> None:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,9 @@ def run(
timeout: int | None = None,
env: t.StrMapping | None = None,
remove_env_keys: t.StrSequence = (),
input_data: str | bytes | None = None,
*,
capture: bool = True,
) -> p.Result[p.Cli.CommandOutput]:
"""Execute a command and require zero exit status."""
...
Expand All @@ -41,6 +44,7 @@ def capture(
timeout: int | None = None,
env: t.StrMapping | None = None,
remove_env_keys: t.StrSequence = (),
input_data: str | bytes | None = None,
) -> p.Result[str]:
"""Execute a command and return stripped stdout."""
...
Expand All @@ -52,7 +56,9 @@ def run_raw(
timeout: int | None = None,
env: t.StrMapping | None = None,
remove_env_keys: t.StrSequence = (),
input_data: bytes | None = None,
input_data: str | bytes | None = None,
*,
capture: bool = True,
) -> p.Result[p.Cli.CommandOutput]:
"""Execute a command without enforcing zero exit status."""
...
Expand All @@ -65,7 +71,7 @@ def run_bytes(
timeout: int | None = None,
env: t.StrMapping | None = None,
remove_env_keys: t.StrSequence = (),
input_data: bytes | None = None,
input_data: str | bytes | None = None,
) -> p.Result[p.Cli.CommandBytesOutput]:
"""Execute a command and preserve byte-exact output."""
...
Expand All @@ -77,10 +83,25 @@ def run_checked(
timeout: int | None = None,
env: t.StrMapping | None = None,
remove_env_keys: t.StrSequence = (),
input_data: str | bytes | None = None,
*,
capture: bool = True,
) -> p.Result[bool]:
"""Execute a command and return a success flag."""
...

def run_live(
self,
cmd: t.StrSequence,
cwd: t.Cli.TextPath | None = None,
timeout: int | None = None,
env: t.StrMapping | None = None,
remove_env_keys: t.StrSequence = (),
input_data: str | bytes | None = None,
) -> p.Result[p.Cli.CommandOutput]:
"""Execute a checked command with inherited live output."""
...

def run_to_file(
self,
cmd: t.StrSequence,
Expand Down
3 changes: 2 additions & 1 deletion src/flext_cli/_utilities/_runtime_commands.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,9 @@ class FlextCliUtilitiesRuntimeCommandsMixin:

if TYPE_CHECKING:

@staticmethod
@classmethod
def run_raw(
cls,
cmd: t.StrSequence,
cwd: t.Cli.TextPath | None = None,
timeout: int | None = None,
Expand Down
8 changes: 4 additions & 4 deletions src/flext_cli/_utilities/_runtime_process_cleanup.py
Original file line number Diff line number Diff line change
Expand Up @@ -76,11 +76,10 @@ def _reap_and_drain(
cls,
process: p.Cli.ProcessHandle,
waiter: threading.Thread,
pump: threading.Thread,
process_done: threading.Event,
wake: threading.Event,
stop: threading.Event,
source: IO[bytes],
pump_streams: tuple[tuple[threading.Thread, IO[bytes]], ...],
cleanup_errors: list[str],
job_handle: int,
absolute_deadline: float | None,
Expand All @@ -98,7 +97,8 @@ def _reap_and_drain(
waiter.join(cls._remaining(cleanup_deadline))
if waiter.is_alive():
cleanup_errors.append("process deadline expired before root reaping")
cls._drain_output(pump, stop, source, cleanup_errors, cleanup_deadline)
for pump, source in pump_streams:
cls._drain_output(pump, stop, source, cleanup_errors, cleanup_deadline)
return return_codes[0] if return_codes else process.poll()

@classmethod
Expand Down Expand Up @@ -157,7 +157,7 @@ def _drain_output(
try:
source.close()
except (OSError, ValueError) as exc:
cleanup_errors.append(f"combined output close error: {exc}")
cleanup_errors.append(f"process output close error: {exc}")
pump.join(cls._remaining(cleanup_deadline))
if pump.is_alive():
cleanup_errors.append("process deadline expired before output drain")
Expand Down
125 changes: 75 additions & 50 deletions src/flext_cli/_utilities/_runtime_process_execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,54 +4,75 @@

import contextlib
import threading
import time
from collections.abc import Callable
from pathlib import Path
from typing import IO, BinaryIO

from flext_cli import c, p, t
from flext_cli import c, p, r, t
from flext_cli._utilities._runtime_process_cleanup import (
FlextCliUtilitiesRuntimeProcessCleanupMixin,
)
from flext_cli._utilities._runtime_process_outcome import (
FlextCliUtilitiesRuntimeProcessOutcomeMixin,
)
from flext_cli._utilities._runtime_process_output import (
FlextCliUtilitiesRuntimeProcessOutputMixin,
)
from flext_cli._utilities._runtime_process_resources import (
FlextCliUtilitiesRuntimeProcessResourcesMixin,
)
from flext_cli._utilities._runtime_process_start import (
FlextCliUtilitiesRuntimeProcessStartMixin,
)
from flext_cli._utilities._runtime_process_timing import (
FlextCliUtilitiesRuntimeProcessTimingMixin,
)


class FlextCliUtilitiesRuntimeProcessExecutionMixin(
FlextCliUtilitiesRuntimeProcessCleanupMixin,
FlextCliUtilitiesRuntimeProcessOutcomeMixin,
FlextCliUtilitiesRuntimeProcessOutputMixin,
FlextCliUtilitiesRuntimeProcessResourcesMixin,
FlextCliUtilitiesRuntimeProcessStartMixin,
FlextCliUtilitiesRuntimeProcessTimingMixin,
):
"""Own one child process and its streaming resources."""

@classmethod
def _execute_streamed_process(
cls,
cmd: t.StrSequence,
output_path: Path,
output_path: Path | None,
cwd: t.Cli.TextPath | None,
env: dict[str, str] | None,
input_data: str | bytes | None,
*,
capture_output: bool,
live: bool,
absolute_deadline: float | None,
grace_seconds: float,
timeout_exit_code: int,
legacy_timeout: bool,
legacy_timeout_seconds: int | None,
) -> p.Result[int]:
timeout: int | None,
deadline: p.Cli.ProcessDeadline | None,
) -> p.Result[p.Cli.CommandBytesOutput]:
"""Own resources and complete one streamed child lifecycle."""
started = time.monotonic()
timing_result = cls._resolve_process_timing(
cmd,
timeout,
deadline,
started,
capture_output=capture_output,
has_output_path=output_path is not None,
live=live,
on_main_thread=threading.current_thread() is threading.main_thread(),
)
if timing_result.failure:
return r[p.Cli.CommandBytesOutput].fail(
timing_result.error or "process deadline resolution failed"
)
absolute_deadline, grace_seconds, timeout_exit_code = timing_result.unwrap()
process: p.Cli.ProcessHandle | None = None
waiter: threading.Thread | None = None
pump: threading.Thread | None = None
source: IO[bytes] | None = None
durable_log: BinaryIO | None = None
job_handle = 0
failures: list[str] = []
Expand All @@ -61,6 +82,9 @@ def _execute_streamed_process(
forwarded_signals: list[int] = []
received_signals: list[int] = []
return_codes: list[int] = []
stdout_output = bytearray()
stderr_output = bytearray()
pump_streams: list[tuple[threading.Thread, IO[bytes]]] = []
pump_stop = threading.Event()
process_done = threading.Event()
wake = threading.Event()
Expand All @@ -77,9 +101,7 @@ def execute_lifecycle() -> None:
final_deadline, \
job_handle, \
process, \
pump, \
return_code, \
source, \
timed_out, \
waiter
if threading.current_thread() is threading.main_thread():
Expand All @@ -92,8 +114,9 @@ def execute_lifecycle() -> None:
if received_signals:
wake.set()
return
output_path.parent.mkdir(parents=True, exist_ok=True)
durable_log = stack.enter_context(output_path.open("wb", buffering=0))
if output_path is not None:
output_path.parent.mkdir(parents=True, exist_ok=True)
durable_log = stack.enter_context(output_path.open("wb", buffering=0))
stdin_result = cls._prepare_streamed_stdin(stack, input_data)
live_result = cls._prepare_live_descriptor(stack, live=live)
if stdin_result.failure:
Expand All @@ -105,32 +128,41 @@ def execute_lifecycle() -> None:
elif cls._spawn_deadline_exhausted(absolute_deadline, grace_seconds):
failures.append("process deadline exhausted before child spawn")
else:
started = cls._start_contained_process(
prepared_cmd, cwd, env, stdin_result.value[0]
combine_output = output_path is not None
pipe_output = combine_output or capture_output
start_result = cls._start_contained_process(
prepared_cmd,
cwd,
env,
stdin_result.value[0],
capture_output=pipe_output,
combine_output=combine_output,
)
if started.failure:
failures.append(started.error or "process start failed")
if start_result.failure:
failures.append(start_result.error or "process start failed")
else:
process, job_handle = started.unwrap()
source = process.stdout
if source is None:
failures.append("process stdout is not available")
return
stack.callback(source.close)
owned_process, job_handle = start_result.unwrap()
process = owned_process
waiter = cls._start_root_waiter(
process, return_codes, failures, process_done, wake
owned_process, return_codes, failures, process_done, wake
)
pump = cls._start_output_pump(
source,
durable_log,
live_result.value[0],
failures,
live_diagnostics,
pump_stop,
wake,
pump_streams.extend(
cls._start_process_output(
owned_process,
stack,
durable_log,
live_result.value[0],
failures,
live_diagnostics,
pump_stop,
wake,
stdout_output,
stderr_output,
capture_output=capture_output,
)
)
timed_out, final_deadline = cls._monitor_process(
process,
owned_process,
process_done,
wake,
failures,
Expand All @@ -140,13 +172,12 @@ def execute_lifecycle() -> None:
grace_seconds,
)
return_code = cls._reap_and_drain(
process,
owned_process,
waiter,
pump,
process_done,
wake,
pump_stop,
source,
tuple(pump_streams),
cleanup_errors,
job_handle,
final_deadline,
Expand All @@ -159,21 +190,14 @@ def execute_lifecycle() -> None:
except c.EXC_OS_VALUE as exc:
failures.append(f"execution error: {exc}")
finally:
if (
process is not None
and waiter is not None
and pump is not None
and source is not None
and not cleanup_complete
):
if process is not None and waiter is not None and not cleanup_complete:
return_code = cls._reap_and_drain(
process,
waiter,
pump,
process_done,
wake,
pump_stop,
source,
tuple(pump_streams),
cleanup_errors,
job_handle,
final_deadline,
Expand All @@ -188,15 +212,16 @@ def execute_lifecycle() -> None:
cleanup_errors.append(close_error)
cleanup_errors.extend(cls._close_process_resources(stack))
cleanup_errors.extend(cls._restore_forwarding_handlers(restore_handlers))
return cls._process_exit_result(
return cls._captured_process_result(
cmd,
return_code,
received_signals,
(*failures, *cleanup_errors),
nonfatal_diagnostics=tuple(live_diagnostics),
stdout_output,
stderr_output,
max(0.0, time.monotonic() - started),
timed_out=timed_out,
legacy_timeout=legacy_timeout,
legacy_timeout_seconds=legacy_timeout_seconds,
timeout_seconds=timeout,
timeout_exit_code=timeout_exit_code,
)

Expand Down
Loading