From 30670b8f83a3b0b23b47b38a31bbac801502e845 Mon Sep 17 00:00:00 2001 From: Viswanath Lekshmanan Date: Thu, 10 Sep 2026 11:40:26 +0530 Subject: [PATCH] feat(cli): add an optional per-job scope to the server job core MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The pre-warmed server runs every job through `_run_command_isolated`, which applies the job's env and cwd, runs the command, and restores the baseline. There was no way for a concern outside this module to bracket that window. Add `register_job_scope_provider`: an optional async context manager entered inside the per-job lock, after env/cwd are applied and before the command runs, so the scope observes the same environment and cwd the job does. It is best-effort — a provider that raises on entry or exit is logged and the job runs unaffected — and its teardown is shielded, so a cancel arriving while the job unwinds cannot leave the scope half torn down. With no provider registered the job path is unchanged. Co-Authored-By: Claude Opus 5 (1M context) --- packages/uipath/pyproject.toml | 2 +- .../uipath/src/uipath/_cli/_server_core.py | 70 ++++- .../uipath/tests/cli/test_server_job_core.py | 258 ++++++++++++++++++ packages/uipath/uv.lock | 2 +- 4 files changed, 327 insertions(+), 5 deletions(-) diff --git a/packages/uipath/pyproject.toml b/packages/uipath/pyproject.toml index 0b1b056fc..86bb4d18f 100644 --- a/packages/uipath/pyproject.toml +++ b/packages/uipath/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "uipath" -version = "2.14.13" +version = "2.14.14" description = "Python SDK and CLI for UiPath Platform, enabling programmatic interaction with automation services, process management, and deployment tools." readme = { file = "README.md", content-type = "text/markdown" } requires-python = ">=3.11" diff --git a/packages/uipath/src/uipath/_cli/_server_core.py b/packages/uipath/src/uipath/_cli/_server_core.py index c426126ff..b884b89c2 100644 --- a/packages/uipath/src/uipath/_cli/_server_core.py +++ b/packages/uipath/src/uipath/_cli/_server_core.py @@ -1,8 +1,11 @@ """Transport-agnostic job core shared by the HTTP and uipath-ipc channels.""" import asyncio +import logging import os import shlex +from collections.abc import AsyncIterator, Callable +from contextlib import AbstractAsyncContextManager, asynccontextmanager from typing import Any from .cli_debug import debug @@ -33,6 +36,66 @@ def init(self) -> None: _state = _ServerState() +logger = logging.getLogger(__name__) + +# Optional per-job scope, registered by an out-of-process concern that needs to bracket +# every job this server runs. ``None`` — the default — means no scope is installed and +# the job runs exactly as it did before. +_job_scope_provider: Callable[[], AbstractAsyncContextManager[None]] | None = None + + +def register_job_scope_provider( + provider: Callable[[], AbstractAsyncContextManager[None]] | None, +) -> None: + """Register an async context manager entered around every job body. + + The scope is entered inside the per-job lock, *after* this module has applied the job's + ``env_vars`` and ``working_dir`` and *before* the command runs, and exited before that + state is restored. A caller therefore sees the same process environment and cwd the job + itself sees, which is what makes per-job setup/teardown possible from outside this module. + + Args: + provider: Zero-argument callable returning a fresh async context manager per job, or + ``None`` to clear the registration. + """ + global _job_scope_provider + _job_scope_provider = provider + + +@asynccontextmanager +async def _job_scope() -> AsyncIterator[None]: + """Enter the registered per-job scope, or do nothing when none is registered. + + Fail-open in both directions: a provider that raises on entry yields an un-scoped job, and + one that raises on exit cannot mask the job's own outcome. Teardown is shielded, so a + cancelled job still runs it to completion. + """ + provider = _job_scope_provider + if provider is None: + yield + return + + scope: AbstractAsyncContextManager[None] | None = None + try: + scope = provider() + await scope.__aenter__() + except Exception: + logger.warning( + "job scope provider failed to start; job runs unscoped", exc_info=True + ) + scope = None + + try: + yield + finally: + if scope is not None: + try: + # Shielded: a cancel delivered while the job unwinds must not leave the + # provider half torn down. + await asyncio.shield(scope.__aexit__(None, None, None)) + except Exception: + logger.warning("job scope provider failed to exit", exc_info=True) + def parse_args(args: str | list[str] | None) -> list[str]: """Parse args into a list of strings.""" @@ -79,9 +142,10 @@ async def _run_command_isolated( "ClientError": True, } - result_value = await asyncio.to_thread( - cmd.main, args, standalone_mode=False - ) + async with _job_scope(): + result_value = await asyncio.to_thread( + cmd.main, args, standalone_mode=False + ) return { "ExitCode": 0, "Error": None, diff --git a/packages/uipath/tests/cli/test_server_job_core.py b/packages/uipath/tests/cli/test_server_job_core.py index 1722d0c5b..d1cbf516b 100644 --- a/packages/uipath/tests/cli/test_server_job_core.py +++ b/packages/uipath/tests/cli/test_server_job_core.py @@ -5,7 +5,10 @@ """ import asyncio +import logging import os +from collections.abc import Iterator +from contextlib import asynccontextmanager from typing import Any from unittest.mock import Mock @@ -92,3 +95,258 @@ def test_parse_args_passes_a_list_through() -> None: def test_parse_args_none_is_empty() -> None: assert _server_core.parse_args(None) == [] + + +# --------------------------------------------------------------------------- +# Per-job scope hook (register_job_scope_provider) +# --------------------------------------------------------------------------- + + +@pytest.fixture +def restore_provider(monkeypatch: pytest.MonkeyPatch) -> None: + """Run the test with no provider registered; restore the original afterwards. + + ``monkeypatch`` records the original value, so it is put back even when the + test registers a provider through :func:`register_job_scope_provider`. + """ + monkeypatch.setattr(_server_core, "_job_scope_provider", None) + + +@pytest.fixture +def capture_core_logs( + caplog: pytest.LogCaptureFixture, +) -> Iterator[pytest.LogCaptureFixture]: + """Attach caplog's handler directly to this module's logger. + + Earlier tests in the suite invoke the Click CLI, whose ``setup_logging`` + leaves the ``uipath`` logger with ``propagate = False``, so records never + reach the root handler caplog normally relies on. + """ + logger = logging.getLogger(_server_core.__name__) + level, propagate = logger.level, logger.propagate + logger.setLevel(logging.DEBUG) + logger.propagate = True + logger.addHandler(caplog.handler) + try: + yield caplog + finally: + logger.removeHandler(caplog.handler) + logger.setLevel(level) + logger.propagate = propagate + + +def _scope_recorder(events: list[tuple[str, str | None, str]]) -> Any: + """A provider whose scope records entry/exit plus the env and cwd it observed.""" + + @asynccontextmanager + async def provider() -> Any: + events.append(("enter", os.environ.get("JOB_MARKER"), os.getcwd())) + try: + yield + finally: + events.append(("exit", os.environ.get("JOB_MARKER"), os.getcwd())) + + return provider + + +async def test_no_provider_registered_runs_the_job_unchanged( + restore_state: Any, restore_provider: Any +) -> None: + _init(restore_state) + cmd = Mock() + cmd.main.return_value = "ok" + result = await _server_core._run_command_isolated(cmd, [], {}, None) + assert result == { + "ExitCode": 0, + "Error": None, + "Result": "ok", + "Unexpected": False, + } + + +async def test_scope_sees_the_jobs_env_and_cwd( + restore_state: Any, restore_provider: Any, tmp_path: Any +) -> None: + """The whole point of the hook: it is entered after env/cwd are applied. + + A caller that had to wrap this function from outside would see neither, which is why the + scope belongs inside. + """ + _init(restore_state) + events: list[tuple[str, str | None, str]] = [] + _server_core.register_job_scope_provider(_scope_recorder(events)) + + cmd = Mock() + cmd.main.return_value = "ok" + await _server_core._run_command_isolated( + cmd, [], {"JOB_MARKER": "job-1"}, str(tmp_path) + ) + + assert [e[0] for e in events] == ["enter", "exit"] + assert {e[1] for e in events} == {"job-1"} # the job's env, not the baseline + assert all( + os.path.realpath(e[2]) == os.path.realpath(str(tmp_path)) for e in events + ) + # ...and the server's own state is restored afterwards. + assert "JOB_MARKER" not in os.environ + + +async def test_scope_exits_before_the_job_env_is_restored( + restore_state: Any, restore_provider: Any +) -> None: + """Teardown inside the scope can still read the job's environment.""" + _init(restore_state) + seen: list[str | None] = [] + + @asynccontextmanager + async def provider() -> Any: + try: + yield + finally: + seen.append(os.environ.get("JOB_MARKER")) + + _server_core.register_job_scope_provider(provider) + cmd = Mock() + cmd.main.return_value = "ok" + await _server_core._run_command_isolated(cmd, [], {"JOB_MARKER": "job-1"}, None) + + assert seen == ["job-1"] + + +async def test_scope_wraps_the_command( + restore_state: Any, restore_provider: Any +) -> None: + """Ordering: enter -> command -> exit.""" + _init(restore_state) + order: list[str] = [] + + @asynccontextmanager + async def provider() -> Any: + order.append("enter") + try: + yield + finally: + order.append("exit") + + _server_core.register_job_scope_provider(provider) + + def record_command(*args: Any, **kwargs: Any) -> str: + order.append("command") + return "ok" + + cmd = Mock() + cmd.main.side_effect = record_command + await _server_core._run_command_isolated(cmd, [], {}, None) + + assert order == ["enter", "command", "exit"] + + +async def test_scope_exits_when_the_command_raises( + restore_state: Any, restore_provider: Any +) -> None: + """The scope must be exited on the job's failure path too.""" + _init(restore_state) + order: list[str] = [] + + @asynccontextmanager + async def provider() -> Any: + order.append("enter") + try: + yield + finally: + order.append("exit") + + _server_core.register_job_scope_provider(provider) + cmd = Mock() + cmd.main.side_effect = RuntimeError("boom") + result = await _server_core._run_command_isolated(cmd, [], {}, None) + + assert order == ["enter", "exit"] + assert result["ExitCode"] == 1 + assert result["Unexpected"] is True + + +async def test_a_provider_that_fails_to_start_does_not_fail_the_job( + restore_state: Any, restore_provider: Any, capture_core_logs: Any +) -> None: + """Best-effort by contract: a broken provider means no scope, not a failed job.""" + _init(restore_state) + + def provider() -> Any: + raise RuntimeError("provider is broken") + + _server_core.register_job_scope_provider(provider) + cmd = Mock() + cmd.main.return_value = "ok" + result = await _server_core._run_command_isolated(cmd, [], {}, None) + + assert result["ExitCode"] == 0 + assert result["Result"] == "ok" + assert "job runs unscoped" in capture_core_logs.text + + +async def test_a_provider_that_fails_to_exit_does_not_change_the_result( + restore_state: Any, restore_provider: Any, capture_core_logs: Any +) -> None: + _init(restore_state) + + @asynccontextmanager + async def provider() -> Any: + yield + raise RuntimeError("exit is broken") + + _server_core.register_job_scope_provider(provider) + cmd = Mock() + cmd.main.return_value = "ok" + result = await _server_core._run_command_isolated(cmd, [], {}, None) + + assert result["ExitCode"] == 0 + assert result["Result"] == "ok" + assert "failed to exit" in capture_core_logs.text + + +async def test_registering_none_removes_the_scope( + restore_state: Any, restore_provider: Any +) -> None: + _init(restore_state) + events: list[tuple[str, str | None, str]] = [] + _server_core.register_job_scope_provider(_scope_recorder(events)) + _server_core.register_job_scope_provider(None) + + cmd = Mock() + cmd.main.return_value = "ok" + await _server_core._run_command_isolated(cmd, [], {}, None) + + assert events == [] + + +async def test_scope_teardown_survives_cancellation( + restore_state: Any, restore_provider: Any +) -> None: + """A cancel arriving during teardown must not leave the provider half torn down.""" + _init(restore_state) + teardown_started = asyncio.Event() + release_teardown = asyncio.Event() + teardown_finished = asyncio.Event() + + @asynccontextmanager + async def provider() -> Any: + try: + yield + finally: + teardown_started.set() + await release_teardown.wait() + teardown_finished.set() + + _server_core.register_job_scope_provider(provider) + cmd = Mock() + cmd.main.return_value = "ok" + + task = asyncio.create_task(_server_core._run_command_isolated(cmd, [], {}, None)) + await asyncio.wait_for(teardown_started.wait(), timeout=5) + task.cancel() + release_teardown.set() + + with pytest.raises(asyncio.CancelledError): + await task + await asyncio.wait_for(teardown_finished.wait(), timeout=5) diff --git a/packages/uipath/uv.lock b/packages/uipath/uv.lock index 18734c70a..d62b3c81c 100644 --- a/packages/uipath/uv.lock +++ b/packages/uipath/uv.lock @@ -2599,7 +2599,7 @@ wheels = [ [[package]] name = "uipath" -version = "2.14.13" +version = "2.14.14" source = { editable = "." } dependencies = [ { name = "applicationinsights" },