Skip to content
Open
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
2 changes: 1 addition & 1 deletion packages/uipath/pyproject.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
70 changes: 67 additions & 3 deletions packages/uipath/src/uipath/_cli/_server_core.py
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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."""
Expand Down Expand Up @@ -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,
Expand Down
258 changes: 258 additions & 0 deletions packages/uipath/tests/cli/test_server_job_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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)
2 changes: 1 addition & 1 deletion packages/uipath/uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading