From 50e0cc7054475b900e6d123137a1f209cbf70164 Mon Sep 17 00:00:00 2001 From: Robert Ursu Date: Thu, 24 Sep 2026 18:26:00 +0300 Subject: [PATCH] feat(cli): really stop a job that uipath server is running [PC-4873] run/debug/eval drive their own event loop on the server's worker thread. They now publish it through run_job_loop to a JobControl, so a stop cancels the job's root task and the runtime unwinds cooperatively, still writing its result. stop_job is shared by IPC StopJob and the new POST /jobs/{key}/stop. It cancels the root task, waits a grace period, cancels every task on the job's loop, and answers False if the job is still running, because it is blocked in a call that only ending the process can interrupt. forceStop shortens the waits. A queued job is dropped before it runs. A stop that targets another resume version, or an unknown job, answers True: that run is not running. The job core no longer hands on the lock, env or cwd while the job thread still runs. A cancelled caller (a dropped IPC connection, a shutdown) stops the job and re-raises only once the thread has exited, and the job scope's teardown completes before the env is restored. A stopped job ends with exit code 143 and "Job stopped on request"; a CancelledError the job raised on its own is reported as an unexpected failure. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01CFMKm7zHGcptS49z4cnnad --- packages/uipath/pyproject.toml | 2 +- .../uipath/src/uipath/_cli/_job_control.py | 99 +++++ .../uipath/src/uipath/_cli/_server_core.py | 274 ++++++++++++-- packages/uipath/src/uipath/_cli/cli_debug.py | 4 +- packages/uipath/src/uipath/_cli/cli_eval.py | 3 +- packages/uipath/src/uipath/_cli/cli_run.py | 4 +- packages/uipath/src/uipath/_cli/cli_server.py | 58 ++- .../uipath/src/uipath/_cli/cli_server_ipc.py | 20 +- .../tests/cli/test_server_cancellation.py | 356 ++++++++++++++++++ packages/uipath/tests/cli/test_server_ipc.py | 23 +- .../uipath/tests/cli/test_server_job_core.py | 15 +- packages/uipath/uv.lock | 2 +- 12 files changed, 795 insertions(+), 65 deletions(-) create mode 100644 packages/uipath/src/uipath/_cli/_job_control.py create mode 100644 packages/uipath/tests/cli/test_server_cancellation.py diff --git a/packages/uipath/pyproject.toml b/packages/uipath/pyproject.toml index 1ce573930..202f9dc63 100644 --- a/packages/uipath/pyproject.toml +++ b/packages/uipath/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "uipath" -version = "2.14.26" +version = "2.14.27" 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/_job_control.py b/packages/uipath/src/uipath/_cli/_job_control.py new file mode 100644 index 000000000..38a017752 --- /dev/null +++ b/packages/uipath/src/uipath/_cli/_job_control.py @@ -0,0 +1,99 @@ +"""Lets the server cancel the job that runs on its worker thread. + +run/debug/eval drive their own event loop inside the ``asyncio.to_thread`` worker. That +loop is the only place a cancellation can land, so the command publishes it here and +the server reaches it through ``loop.call_soon_threadsafe``. + +Kept import-light: ``_cli/__init__.py`` defers heavy imports. +""" + +import asyncio +import contextvars +import threading +from typing import Any + +CURRENT_JOB_CONTROL: contextvars.ContextVar["JobControl | None"] = ( + contextvars.ContextVar("uipath_current_job_control", default=None) +) + + +class JobControl: + """Handle on one job's event loop, shared by the server loop and the job thread.""" + + def __init__(self) -> None: + """Create an unbound control; the job binds its loop once it has one.""" + self.cancel_requested = False + self._delivered = False + self._loop: asyncio.AbstractEventLoop | None = None + self._task: "asyncio.Task[Any] | None" = None + self._sync = threading.Lock() + + def bind(self, loop: asyncio.AbstractEventLoop, task: "asyncio.Task[Any]") -> None: + """Publish the job's loop and root task (job thread).""" + with self._sync: + self._loop, self._task = loop, task + pending = self.cancel_requested + if pending: + self.cancel() + + def unbind(self) -> None: + """Withdraw the loop once the root task is done (job thread).""" + with self._sync: + self._loop = self._task = None + + def cancel(self) -> None: + """Cancel the job's root task, at most once. + + A second delivery would land inside the runtime's cleanup ``finally`` blocks, + the ones that write ``output.json``, and abort them. Before the job has a loop, + the request is recorded and applied by ``bind``. + """ + with self._sync: + self.cancel_requested = True + if self._delivered or self._loop is None or self._task is None: + return + loop, task = self._loop, self._task + self._delivered = True + try: + loop.call_soon_threadsafe(task.cancel) + except RuntimeError: + with self._sync: + self._delivered = False + + def cancel_all(self) -> None: + """Cancel every task on the job's loop, giving up on a clean cleanup.""" + with self._sync: + self.cancel_requested = True + loop = self._loop + if loop is None: + return + + def _sweep() -> None: + for task in asyncio.all_tasks(loop): + task.cancel() + + try: + loop.call_soon_threadsafe(_sweep) + except RuntimeError: + pass + + +def run_job_loop(coro: Any) -> Any: + """``asyncio.run`` that publishes its loop and root task to the job's control. + + With no control in scope (``uipath run`` from a terminal) this is ``asyncio.run``. + The loop is withdrawn before the runner closes, so a sweep can never cancel the + runner's own wait on the job's executor threads. + """ + control = CURRENT_JOB_CONTROL.get() + if control is None: + return asyncio.run(coro) + + with asyncio.Runner() as runner: + loop = runner.get_loop() + task = loop.create_task(coro, context=contextvars.copy_context()) + control.bind(loop, task) + try: + return loop.run_until_complete(task) + finally: + control.unbind() diff --git a/packages/uipath/src/uipath/_cli/_server_core.py b/packages/uipath/src/uipath/_cli/_server_core.py index f04ae1b5f..5a0241d47 100644 --- a/packages/uipath/src/uipath/_cli/_server_core.py +++ b/packages/uipath/src/uipath/_cli/_server_core.py @@ -1,13 +1,17 @@ """Transport-agnostic job core shared by the HTTP and uipath-ipc channels.""" import asyncio +import contextvars +import functools import logging import os import shlex from collections.abc import AsyncIterator, Callable from contextlib import AbstractAsyncContextManager, asynccontextmanager +from dataclasses import dataclass, field from typing import Any +from ._job_control import CURRENT_JOB_CONTROL, JobControl from .cli_debug import debug from .cli_eval import eval from .cli_run import run @@ -19,12 +23,33 @@ } +# Distinct from any click exit code, so a stop on request is not read as a failure. +EXIT_CODE_STOPPED = 143 + +STOP_GRACE_SECONDS = 30.0 +STOP_ESCALATION_SECONDS = 10.0 +FORCE_STOP_GRACE_SECONDS = 5.0 +FORCE_STOP_ESCALATION_SECONDS = 5.0 + + +@dataclass +class _Job: + job_key: str + resume_version: int | None + task: "asyncio.Task[Any] | None" + control: JobControl = field(default_factory=JobControl) + started: bool = False + finished: asyncio.Event = field(default_factory=asyncio.Event) + + class _ServerState: """Mutable server state, initialized lazily at server startup.""" def __init__(self) -> None: self.lock: asyncio.Lock | None = None self.baseline_env: dict[str, str] | None = None + # Queued and running jobs; the lock lets at most one of them run. + self.jobs: list[_Job] = [] # Optional per-job scope; ``None`` (the default) leaves jobs unscoped. self.job_scope_provider: ( Callable[[], AbstractAsyncContextManager[None]] | None @@ -54,10 +79,8 @@ def register_job_scope_provider( and cwd the job itself sees, which is what makes per-job setup/teardown possible from outside this module. - Teardown ordering is best-effort rather than guaranteed. When a job is cancelled the - teardown is shielded but not awaited, so it runs detached and may observe the restored - server baseline instead of the job's env and cwd. A provider that needs the job's values - on the teardown side should capture them on entry rather than read the process state. + The job's env and cwd are restored only after the teardown has finished, cancelled or + not. Args: provider: Zero-argument callable returning a fresh async context manager per job, or @@ -71,8 +94,8 @@ 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. + one that raises on exit cannot mask the job's own outcome. A cancelled job still runs the + teardown to completion before the cancellation moves on. """ provider = _state.job_scope_provider if provider is None: @@ -93,12 +116,33 @@ async def _job_scope() -> AsyncIterator[None]: 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) + teardown = asyncio.ensure_future(scope.__aexit__(None, None, None)) + interrupted = await _await_despite_cancellation(teardown) + if not teardown.cancelled() and teardown.exception() is not None: + logger.warning( + "job scope provider failed to exit", exc_info=teardown.exception() + ) + if interrupted: + raise asyncio.CancelledError + + +async def _await_despite_cancellation( + future: "asyncio.Future[Any]", on_cancel: Callable[[], None] | None = None +) -> bool: + """Wait for ``future`` to finish even if this task is cancelled meanwhile. + + Returns whether a cancellation arrived, so the caller can re-raise it once its own + cleanup is done. ``on_cancel`` runs on each cancellation that arrives while waiting. + """ + interrupted = False + while not future.done(): + try: + await asyncio.wait([future]) + except asyncio.CancelledError: + interrupted = True + if on_cancel is not None: + on_cancel() + return interrupted def parse_args(args: str | list[str] | None) -> list[str]: @@ -121,6 +165,154 @@ def _exit_code_outcome(exit_code: int) -> dict[str, Any]: } +def _stopped_outcome() -> dict[str, Any]: + return { + "ExitCode": EXIT_CODE_STOPPED, + "Error": "Job stopped on request", + "Result": None, + "Unexpected": False, + "Stopped": True, + } + + +def _find_job(job_key: str, resume_version: int | None) -> _Job | None: + for job in _state.jobs: + if job.job_key != job_key: + continue + if ( + resume_version is not None + and job.resume_version is not None + and job.resume_version != resume_version + ): + continue + return job + return None + + +async def _wait_finished(job: _Job, timeout: float) -> bool: + try: + await asyncio.wait_for(job.finished.wait(), timeout) + return True + except asyncio.TimeoutError: + return False + + +async def stop_job( + job_key: str, resume_version: int | None = None, force: bool = False +) -> bool: + """Stop a queued or running job and report whether it is no longer running. + + A running job is stopped cooperatively: its event loop's root task is cancelled, so + the runtime unwinds and still writes its result. If it has not finished within the + grace period every task on its loop is cancelled, and if it still has not finished + the answer is False: it is blocked in a call that cannot be interrupted, and only + ending the process stops it. ``force`` shortens both waits. + + A job that is not known, or whose live run has another resume version, is not + running, so the answer is True. + """ + job = _find_job(job_key, resume_version) + if job is None: + logger.info("StopJob for %s: no such job is queued or running", job_key) + return True + + grace, escalation = ( + (FORCE_STOP_GRACE_SECONDS, FORCE_STOP_ESCALATION_SECONDS) + if force + else (STOP_GRACE_SECONDS, STOP_ESCALATION_SECONDS) + ) + + job.control.cancel() + if not job.started: + if job.task is not None: + job.task.cancel() + return True + + if await _wait_finished(job, grace): + return True + + logger.warning( + "StopJob for %s: still running after %ss; cancelling every task on its loop", + job_key, + grace, + ) + job.control.cancel_all() + if await _wait_finished(job, escalation): + return True + + logger.error( + "StopJob for %s: still running; it is blocked in a call that cannot be " + "interrupted", + job_key, + ) + return False + + +async def _run_job_thread(cmd: Any, args: list[str], job: _Job) -> dict[str, Any]: + """Run the command on a worker thread and return once that thread has exited. + + Cancelling the awaiting task (a dropped connection, a server shutdown) stops the job + instead of abandoning it, and the cancellation is re-raised only once the thread is + gone, so the lock, env and cwd are never handed on while the job still runs. + """ + loop = asyncio.get_running_loop() + token = CURRENT_JOB_CONTROL.set(job.control) + try: + context = contextvars.copy_context() + finally: + CURRENT_JOB_CONTROL.reset(token) + # A plain executor future rather than a task: a task re-raises the job's SystemExit + # into the server loop instead of keeping it as the job's outcome. + worker = loop.run_in_executor( + None, functools.partial(context.run, cmd.main, args, standalone_mode=False) + ) + escalations: list[asyncio.TimerHandle] = [] + + def _stop_abandoned_job() -> None: + job.control.cancel() + if not escalations: + escalations.append( + loop.call_later(STOP_GRACE_SECONDS, job.control.cancel_all) + ) + + try: + interrupted = await _await_despite_cancellation(worker, _stop_abandoned_job) + finally: + for handle in escalations: + handle.cancel() + job.finished.set() + + if interrupted: + raise asyncio.CancelledError + + if worker.cancelled() or isinstance(worker.exception(), asyncio.CancelledError): + if job.control.cancel_requested: + return _stopped_outcome() + return { + "ExitCode": 1, + "Error": "Job cancelled itself", + "Result": None, + "Unexpected": True, + } + + exc = worker.exception() + if exc is not None: + raise exc + + result_value = worker.result() + # Under standalone_mode=False click returns ctx.exit(N)'s code instead of raising, + # and every ConsoleLogger.error path ends in ctx.exit(1). run/debug/eval never + # return an int of their own, so an int here is always an exit code. + if isinstance(result_value, int) and not isinstance(result_value, bool): + return _exit_code_outcome(result_value) + return { + "ExitCode": 0, + "Error": None, + "Result": result_value, + "Unexpected": False, + } + + async def _run_command_isolated( cmd: Any, args: list[str], @@ -128,17 +320,50 @@ async def _run_command_isolated( working_dir: str | None, on_run_start: Callable[[], None] | None = None, on_run_end: Callable[[], None] | None = None, + job_key: str | None = None, + resume_version: int | None = None, ) -> dict[str, Any]: """Run one command with per-job env/cwd isolation (the shared job core). - ``on_run_start`` / ``on_run_end`` run inside the serialization lock. That orders them against - the next job's start, but claims nothing once this task is cancelled: the job runs on a thread - that cancellation cannot reach, so it outlives the lock and the globals move under it. + ``on_run_start`` / ``on_run_end`` run inside the serialization lock, and the lock, + env and cwd are only given back once the job's thread has exited. A job with a + ``job_key`` can be stopped through :func:`stop_job`. """ if _state.lock is None or _state.baseline_env is None: raise RuntimeError("Server state not initialized") - async with _state.lock: + job = _Job(job_key or "", resume_version, asyncio.current_task()) + if job_key: + _state.jobs.append(job) + try: + return await _run_registered_job( + job, cmd, args, env_vars, working_dir, on_run_start, on_run_end + ) + finally: + if job_key: + _state.jobs.remove(job) + + +async def _run_registered_job( + job: _Job, + cmd: Any, + args: list[str], + env_vars: dict[str, str], + working_dir: str | None, + on_run_start: Callable[[], None] | None, + on_run_end: Callable[[], None] | None, +) -> dict[str, Any]: + assert _state.lock is not None and _state.baseline_env is not None + try: + await _state.lock.acquire() + except asyncio.CancelledError: + if job.control.cancel_requested: + asyncio.current_task().uncancel() # type: ignore[union-attr] + return _stopped_outcome() + raise + job.started = True + + try: original_cwd = os.getcwd() try: # Start from server baseline + request env vars only, so nothing from @@ -169,23 +394,10 @@ async def _run_command_isolated( # entered and cleared after it exits, so anything the provider logs on the # way in or out belongs to this job rather than to the server. async with _job_scope(): - result_value = await asyncio.to_thread( - cmd.main, args, standalone_mode=False - ) + return await _run_job_thread(cmd, args, job) finally: if on_run_end is not None: on_run_end() - # Under standalone_mode=False click returns ctx.exit(N)'s code instead of raising, - # and every ConsoleLogger.error path ends in ctx.exit(1). run/debug/eval never - # return an int of their own, so an int here is always an exit code. - if isinstance(result_value, int) and not isinstance(result_value, bool): - return _exit_code_outcome(result_value) - return { - "ExitCode": 0, - "Error": None, - "Result": result_value, - "Unexpected": False, - } except SystemExit as e: return _exit_code_outcome(e.code if isinstance(e.code, int) else 1) except Exception as e: # report any job failure as a result, not a fault @@ -198,3 +410,5 @@ async def _run_command_isolated( pass os.environ.clear() os.environ.update(_state.baseline_env) + finally: + _state.lock.release() diff --git a/packages/uipath/src/uipath/_cli/cli_debug.py b/packages/uipath/src/uipath/_cli/cli_debug.py index fc2372f0d..aecfeca9a 100644 --- a/packages/uipath/src/uipath/_cli/cli_debug.py +++ b/packages/uipath/src/uipath/_cli/cli_debug.py @@ -1,4 +1,3 @@ -import asyncio import logging from typing import Any, cast, get_args @@ -29,6 +28,7 @@ from uipath.tracing import LiveTrackingSpanProcessor, LlmOpsHttpExporter from ._governance_bootstrap import GovernanceBootstrap, resolve_governance +from ._job_control import run_job_loop from ._run_telemetry import RunTelemetry from ._telemetry import track_command from ._utils._console import ConsoleLogger @@ -311,7 +311,7 @@ async def execute_debug_runtime(): finally: trace_manager.shutdown() - asyncio.run(execute_debug_runtime()) + run_job_loop(execute_debug_runtime()) except Exception as e: console.error( f"Error occurred: {e or 'Execution failed'}", include_traceback=True diff --git a/packages/uipath/src/uipath/_cli/cli_eval.py b/packages/uipath/src/uipath/_cli/cli_eval.py index 2a446c95a..4f01ee1aa 100644 --- a/packages/uipath/src/uipath/_cli/cli_eval.py +++ b/packages/uipath/src/uipath/_cli/cli_eval.py @@ -39,6 +39,7 @@ LlmOpsHttpExporter, ) +from ._job_control import run_job_loop from ._utils._console import ConsoleLogger logger = logging.getLogger(__name__) @@ -528,7 +529,7 @@ async def execute_eval(): finally: await runtime_factory.dispose() - asyncio.run(execute_eval()) + run_job_loop(execute_eval()) except _EvalDiscoveryError as e: click.echo("\n".join(e.get_usage_help())) diff --git a/packages/uipath/src/uipath/_cli/cli_run.py b/packages/uipath/src/uipath/_cli/cli_run.py index 774b3dba6..b5ca683e0 100644 --- a/packages/uipath/src/uipath/_cli/cli_run.py +++ b/packages/uipath/src/uipath/_cli/cli_run.py @@ -1,4 +1,3 @@ -import asyncio from typing import Any import click @@ -36,6 +35,7 @@ from ._errors import EntrypointDiscoveryException from ._governance_bootstrap import GovernanceBootstrap, resolve_governance +from ._job_control import run_job_loop from ._run_telemetry import RunTelemetry from ._telemetry import track_command from ._utils._console import ConsoleLogger @@ -367,7 +367,7 @@ async def execute() -> None: finally: trace_manager.shutdown() - asyncio.run(execute()) + run_job_loop(execute()) except _RunDiscoveryError as e: click.echo("\n".join(e.get_usage_help())) diff --git a/packages/uipath/src/uipath/_cli/cli_server.py b/packages/uipath/src/uipath/_cli/cli_server.py index c6d6bf71a..41c6ec535 100644 --- a/packages/uipath/src/uipath/_cli/cli_server.py +++ b/packages/uipath/src/uipath/_cli/cli_server.py @@ -17,6 +17,7 @@ _run_command_isolated, _state, parse_args, + stop_job, ) from ._telemetry import track_command from ._utils._console import ConsoleLogger @@ -184,9 +185,26 @@ async def handle_start(request: web.Request) -> web.Response: status=400, ) + resume_version = get_field(message, "resumeVersion", "ResumeVersion") + if resume_version is not None and not isinstance(resume_version, int): + return web.json_response( + { + "success": False, + "error": "Invalid field: 'resumeVersion' must be an int", + }, + status=400, + ) + console.info(f"Starting job {job_key}: {command_name} {args}") - result = await _run_command_isolated(cmd, args, env_vars, working_dir) + result = await _run_command_isolated( + cmd, + args, + env_vars, + working_dir, + job_key=job_key, + resume_version=resume_version, + ) # The .NET peer decides success from ``exitCode`` alone and defaults a missing one to 0. exit_code = result["ExitCode"] @@ -230,6 +248,43 @@ async def handle_start(request: web.Request) -> web.Response: ) +async def handle_stop(request: web.Request) -> web.Response: + """Handle POST /jobs/{job_key}/stop — 200 with whether the job no longer runs.""" + job_key = request.match_info.get("job_key") + if not job_key: + return web.json_response( + {"success": False, "error": "Missing job_key"}, status=400 + ) + + message: dict[str, Any] = {} + if request.can_read_body: + try: + message = await request.json() + except json.JSONDecodeError: + return web.json_response( + {"success": False, "error": "Invalid JSON"}, status=400 + ) + if not isinstance(message, dict): + return web.json_response( + {"success": False, "error": "Invalid JSON"}, status=400 + ) + + resume_version = get_field(message, "resumeVersion", "ResumeVersion") + if resume_version is not None and not isinstance(resume_version, int): + return web.json_response( + { + "success": False, + "error": "Invalid field: 'resumeVersion' must be an int", + }, + status=400, + ) + force = get_field(message, "forceStop", "ForceStop") is True + + console.info(f"StopJob requested for {job_key} (force={force})") + stopped = await stop_job(job_key, resume_version, force=force) + return web.json_response({"success": True, "job_key": job_key, "stopped": stopped}) + + ALLOWED_HOSTS = {"127.0.0.1", "localhost", "[::1]"} @@ -264,6 +319,7 @@ def create_app() -> web.Application: app = web.Application(middlewares=[host_validation_middleware]) app.router.add_get("/health", handle_health) app.router.add_post("/jobs/{job_key}/start", handle_start) + app.router.add_post("/jobs/{job_key}/stop", handle_stop) return app diff --git a/packages/uipath/src/uipath/_cli/cli_server_ipc.py b/packages/uipath/src/uipath/_cli/cli_server_ipc.py index b329159e9..bc01e3b26 100644 --- a/packages/uipath/src/uipath/_cli/cli_server_ipc.py +++ b/packages/uipath/src/uipath/_cli/cli_server_ipc.py @@ -3,7 +3,13 @@ from dataclasses import dataclass, field from typing import TYPE_CHECKING, Any -from ._server_core import COMMANDS, _run_command_isolated, _state, parse_args +from ._server_core import ( + COMMANDS, + _run_command_isolated, + _state, + parse_args, + stop_job, +) from ._utils._console import ConsoleLogger if TYPE_CHECKING: @@ -69,7 +75,7 @@ async def RunJob( @abstractmethod async def StopJob(self, request: PythonServerStopJobRequest) -> bool: - """Cancel a running job by key (bool return avoids fire-and-forget).""" + """Stop a job; True once it no longer runs (bool return avoids fire-and-forget).""" class PythonRuntimeService(IPythonRuntimeServer): @@ -147,6 +153,8 @@ def _install() -> None: request.workingDirectory, on_run_start=on_run_start, on_run_end=on_run_end, + job_key=request.jobKey or None, + resume_version=request.resumeVersion, ) # Await, don't block: these share this loop, and the peer may unregister the job once this returns. @@ -160,11 +168,13 @@ def _install() -> None: ) async def StopJob(self, request: PythonServerStopJobRequest) -> bool: - console.warning( + console.info( f"StopJob requested for {_run_id(request.jobKey, request.resumeVersion)} " - f"(force={request.forceStop}), but this server cannot stop a running job" + f"(force={request.forceStop})" + ) + return await stop_job( + request.jobKey, request.resumeVersion, force=request.forceStop ) - return False async def start_ipc_server(pipe_name: str) -> None: diff --git a/packages/uipath/tests/cli/test_server_cancellation.py b/packages/uipath/tests/cli/test_server_cancellation.py new file mode 100644 index 000000000..161f46454 --- /dev/null +++ b/packages/uipath/tests/cli/test_server_cancellation.py @@ -0,0 +1,356 @@ +"""Stopping a job that ``uipath server`` is running. + +The commands below are shaped like run/debug/eval: a click command whose body ends in +``run_job_loop(...)`` on the worker thread. +""" + +import asyncio +import os +import threading +import time +from typing import Any + +import click +import pytest +from aiohttp.test_utils import TestClient, TestServer + +from uipath._cli import _server_core, cli_server, cli_server_ipc +from uipath._cli._job_control import CURRENT_JOB_CONTROL, run_job_loop +from uipath._cli._server_core import ( + EXIT_CODE_STOPPED, + _run_command_isolated, + _ServerState, + stop_job, +) + +JOB = "3f2504e0-4f89-11d3-9a0c-0305e82c3301" +OTHER_JOB = "5b1f6c9e-2d4a-4e8b-9f3c-7a6d5e4c3b2a" + +started = threading.Event() +cleanup_ran = threading.Event() +release = threading.Event() + + +@pytest.fixture(autouse=True) +def fresh_state(monkeypatch: pytest.MonkeyPatch) -> _ServerState: + state = _ServerState() + state.lock = asyncio.Lock() + state.baseline_env = dict(os.environ) + monkeypatch.setattr(_server_core, "_state", state) + monkeypatch.setattr(cli_server_ipc, "_state", state) + started.clear() + cleanup_ran.clear() + release.clear() + return state + + +@pytest.fixture +def short_waits(monkeypatch: pytest.MonkeyPatch) -> None: + for name in ( + "STOP_GRACE_SECONDS", + "STOP_ESCALATION_SECONDS", + "FORCE_STOP_GRACE_SECONDS", + "FORCE_STOP_ESCALATION_SECONDS", + ): + monkeypatch.setattr(_server_core, name, 0.2) + + +@click.command() +def long_job() -> None: + async def body() -> None: + try: + started.set() + await asyncio.sleep(30) + finally: + cleanup_ran.set() + + run_job_loop(body()) + + +@click.command() +def slow_cleanup_job() -> None: + async def body() -> None: + try: + started.set() + await asyncio.sleep(30) + finally: + await asyncio.sleep(0.3) + cleanup_ran.set() + + run_job_loop(body()) + + +@click.command() +def late_loop_job() -> None: + started.set() + release.wait(10) + + async def body() -> None: + try: + await asyncio.sleep(30) + finally: + cleanup_ran.set() + + run_job_loop(body()) + + +@click.command() +def nested_thread_job() -> None: + async def body() -> None: + started.set() + await asyncio.to_thread(time.sleep, 0.5) + + run_job_loop(body()) + + +@click.command() +def self_cancelling_job() -> None: + async def body() -> None: + started.set() + raise asyncio.CancelledError() + + run_job_loop(body()) + + +@click.command() +def blocking_job() -> None: + started.set() + release.wait(10) + + +@click.command() +def quick_job() -> None: + async def body() -> None: + await asyncio.sleep(0) + + run_job_loop(body()) + + +async def wait_until(event: threading.Event, timeout: float = 10.0) -> None: + deadline = time.monotonic() + timeout + while not event.is_set(): + assert time.monotonic() < deadline, "timed out" + await asyncio.sleep(0.01) + + +def start(cmd: Any, job_key: str = JOB, **kwargs: Any) -> "asyncio.Task[Any]": + return asyncio.create_task( + _run_command_isolated(cmd, [], {}, None, job_key=job_key, **kwargs) + ) + + +def test_run_job_loop_is_asyncio_run_outside_the_server() -> None: + assert CURRENT_JOB_CONTROL.get() is None + + async def body() -> int: + return 42 + + assert run_job_loop(body()) == 42 + + +async def test_stop_unwinds_a_running_job_and_frees_the_lock( + fresh_state: _ServerState, +) -> None: + job = start(long_job) + await wait_until(started) + + assert await stop_job(JOB) is True + outcome = await job + + assert outcome["ExitCode"] == EXIT_CODE_STOPPED + assert outcome["Stopped"] is True + assert cleanup_ran.is_set() + assert fresh_state.lock is not None and not fresh_state.lock.locked() + assert fresh_state.jobs == [] + + +async def test_a_stop_before_the_job_has_a_loop_is_applied_when_it_gets_one() -> None: + job = start(late_loop_job) + await wait_until(started) + + stopping = asyncio.create_task(stop_job(JOB)) + await asyncio.sleep(0.05) + release.set() + + assert await stopping is True + assert (await job)["ExitCode"] == EXIT_CODE_STOPPED + assert cleanup_ran.is_set() + + +async def test_a_repeated_stop_does_not_abort_the_cleanup() -> None: + job = start(slow_cleanup_job) + await wait_until(started) + + first = asyncio.create_task(stop_job(JOB)) + await asyncio.sleep(0.05) + second = asyncio.create_task(stop_job(JOB, force=True)) + + assert await first is True + assert await second is True + assert (await job)["ExitCode"] == EXIT_CODE_STOPPED + assert cleanup_ran.is_set() + + +async def test_a_job_waiting_on_its_own_thread_stops_once_that_call_returns() -> None: + job = start(nested_thread_job) + await wait_until(started) + + assert await stop_job(JOB) is True + assert (await job)["ExitCode"] == EXIT_CODE_STOPPED + + +async def test_an_uninterruptible_job_reports_false_and_keeps_the_lock( + fresh_state: _ServerState, short_waits: None +) -> None: + job = start(blocking_job) + await wait_until(started) + + assert await stop_job(JOB) is False + assert fresh_state.lock is not None and fresh_state.lock.locked() + assert not job.done() + + release.set() + outcome = await job + assert outcome["ExitCode"] == 0 + assert not fresh_state.lock.locked() + + +async def test_a_self_inflicted_cancellation_is_a_fault_not_a_stop() -> None: + outcome = await start(self_cancelling_job) + + assert outcome["ExitCode"] == 1 + assert outcome["Unexpected"] is True + assert "Stopped" not in outcome + + +async def test_a_queued_job_is_stopped_without_running( + fresh_state: _ServerState, +) -> None: + running = start(long_job) + await wait_until(started) + queued = start(quick_job, job_key=OTHER_JOB) + await asyncio.sleep(0.05) + + assert await stop_job(OTHER_JOB) is True + assert (await queued)["ExitCode"] == EXIT_CODE_STOPPED + assert not running.done() + + assert await stop_job(JOB) is True + await running + assert fresh_state.jobs == [] + + +async def test_a_stop_for_another_resume_version_leaves_the_live_run_alone() -> None: + job = start(long_job, resume_version=2) + await wait_until(started) + + assert await stop_job(JOB, resume_version=1) is True + await asyncio.sleep(0.05) + assert not job.done() + assert not cleanup_ran.is_set() + + assert await stop_job(JOB, resume_version=2) is True + assert (await job)["ExitCode"] == EXIT_CODE_STOPPED + + +async def test_stopping_an_unknown_job_reports_it_is_not_running() -> None: + assert await stop_job("no-such-job") is True + + +async def test_a_cancelled_caller_stops_the_job_and_waits_for_its_thread( + fresh_state: _ServerState, +) -> None: + job = asyncio.create_task( + _run_command_isolated( + slow_cleanup_job, [], {"JOB_MARKER": "job-1"}, None, job_key=JOB + ) + ) + await wait_until(started) + + job.cancel() + with pytest.raises(asyncio.CancelledError): + await job + + assert cleanup_ran.is_set() + assert fresh_state.lock is not None and not fresh_state.lock.locked() + assert os.environ.get("JOB_MARKER") is None + + +async def test_a_cancelled_caller_keeps_the_lock_while_the_thread_runs( + fresh_state: _ServerState, +) -> None: + job = asyncio.create_task( + _run_command_isolated( + blocking_job, [], {"JOB_MARKER": "job-1"}, None, job_key=JOB + ) + ) + await wait_until(started) + + job.cancel() + await asyncio.sleep(0.1) + assert not job.done() + assert fresh_state.lock is not None and fresh_state.lock.locked() + assert os.environ.get("JOB_MARKER") == "job-1" + + release.set() + with pytest.raises(asyncio.CancelledError): + await job + assert not fresh_state.lock.locked() + assert os.environ.get("JOB_MARKER") is None + + +async def test_ipc_stop_job_stops_a_running_run_job( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setitem(_server_core.COMMANDS, "long", long_job) + service = cli_server_ipc.PythonRuntimeService() + run = asyncio.create_task( + service.RunJob( + cli_server_ipc.PythonServerRunRequest( + jobKey=JOB, resumeVersion=3, command="long" + ) + ) + ) + await wait_until(started) + + stopped = await service.StopJob( + cli_server_ipc.PythonServerStopJobRequest(jobKey=JOB, resumeVersion=3) + ) + + assert stopped is True + result = await run + assert result.exitCode == EXIT_CODE_STOPPED + assert result.error == "Job stopped on request" + + +async def test_http_stop_route_stops_a_running_job( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setitem(_server_core.COMMANDS, "long", long_job) + async with TestClient(TestServer(cli_server.create_app())) as client: + start_request = asyncio.create_task( + client.post(f"/jobs/{JOB}/start", json={"command": "long"}) + ) + await wait_until(started) + + stop = await client.post(f"/jobs/{JOB}/stop", json={"forceStop": True}) + assert stop.status == 200 + assert (await stop.json())["stopped"] is True + + response = await start_request + body = await response.json() + assert body["success"] is False + assert body["exitCode"] == EXIT_CODE_STOPPED + + +async def test_http_stop_route_accepts_an_empty_body() -> None: + async with TestClient(TestServer(cli_server.create_app())) as client: + response = await client.post(f"/jobs/{JOB}/stop") + assert response.status == 200 + assert (await response.json())["stopped"] is True + + +async def test_http_stop_route_rejects_a_bad_resume_version() -> None: + async with TestClient(TestServer(cli_server.create_app())) as client: + response = await client.post(f"/jobs/{JOB}/stop", json={"resumeVersion": "two"}) + assert response.status == 400 diff --git a/packages/uipath/tests/cli/test_server_ipc.py b/packages/uipath/tests/cli/test_server_ipc.py index b3b4a9c69..18798d44a 100644 --- a/packages/uipath/tests/cli/test_server_ipc.py +++ b/packages/uipath/tests/cli/test_server_ipc.py @@ -189,16 +189,15 @@ def test_stop_job_accepts_resume_version_and_force_stop(self, pipe): } result = asyncio.run(_with_proxy(pipe, lambda p: p.StopJob(request))) - assert result is False + assert result is True - def test_stop_job_reports_that_nothing_was_stopped(self, pipe): - """This server cannot stop a running job, so it must not claim it did.""" + def test_stopping_an_unknown_job_reports_it_is_not_running(self, pipe): result = asyncio.run( _with_proxy( pipe, lambda p: p.StopJob({"jobKey": "job-1", "forceStop": True}) ) ) - assert result is False + assert result is True class TestIpcServerEnvIsolation: @@ -463,7 +462,9 @@ def _record( monkeypatch.setattr(_job_api, "install_runtime_sinks", _record) - async def _fake_run(cmd, args, env, wd, on_run_start=None, on_run_end=None): + async def _fake_run( + cmd, args, env, wd, on_run_start=None, on_run_end=None, **_: Any + ): if on_run_start: on_run_start() # Emit through the installed sink, so a log line actually crosses the pipe. @@ -541,7 +542,9 @@ def test_runjob_installs_the_sinks_from_the_request_callback(self, monkeypatch): _job_api, "clear_runtime_sinks", lambda: events.append(("clear",)) ) - async def _fake_run(cmd, args, env, wd, on_run_start=None, on_run_end=None): + async def _fake_run( + cmd, args, env, wd, on_run_start=None, on_run_end=None, **_: Any + ): if on_run_start: on_run_start() events.append(("run",)) @@ -581,7 +584,9 @@ async def aflush_pending(self, *a: Any, **k: Any) -> None: ) monkeypatch.setattr(_job_api, "clear_runtime_sinks", lambda: None) - async def _fake_run(cmd, args, env, wd, on_run_start=None, on_run_end=None): + async def _fake_run( + cmd, args, env, wd, on_run_start=None, on_run_end=None, **_: Any + ): if on_run_start: on_run_start() events.append("run") @@ -620,7 +625,9 @@ def test_runjob_skips_sinks_when_not_opted_in(self, monkeypatch): _job_api, "clear_runtime_sinks", lambda: events.append(("clear",)) ) - async def _fake_run(cmd, args, env, wd, on_run_start=None, on_run_end=None): + async def _fake_run( + cmd, args, env, wd, on_run_start=None, on_run_end=None, **_: Any + ): if on_run_start: on_run_start() events.append(("run",)) diff --git a/packages/uipath/tests/cli/test_server_job_core.py b/packages/uipath/tests/cli/test_server_job_core.py index 652775455..67e764613 100644 --- a/packages/uipath/tests/cli/test_server_job_core.py +++ b/packages/uipath/tests/cli/test_server_job_core.py @@ -460,23 +460,10 @@ async def provider() -> Any: await asyncio.wait_for(teardown_finished.wait(), timeout=5) -@pytest.mark.xfail( - strict=True, - reason="Teardown is shielded but not awaited, so under cancellation it runs detached " - "and observes the restored baseline instead of the job's env. Open review thread on " - "the job-scope PR; remove this marker with the fix.", -) async def test_scope_teardown_sees_the_job_env_when_cancelled( restore_state: Any, restore_provider: Any ) -> None: - """Pin the ordering half of the scope contract, which cancellation currently breaks. - - On the normal path the scope exits before the job's env is restored, which is what lets a - provider do per-job teardown against the job's own values. Under cancellation that - ordering is documented as best-effort and does not hold. This asserts the behaviour worth - having, so the gap stays visible and flips loudly if it is ever closed — no other test in - this file can observe the difference. - """ + """The scope exits before the job's env is restored, even when the job is cancelled.""" _init(restore_state) teardown_started = asyncio.Event() observed: list[str | None] = [] diff --git a/packages/uipath/uv.lock b/packages/uipath/uv.lock index ebf95dbc7..863571f7c 100644 --- a/packages/uipath/uv.lock +++ b/packages/uipath/uv.lock @@ -2599,7 +2599,7 @@ wheels = [ [[package]] name = "uipath" -version = "2.14.26" +version = "2.14.27" source = { editable = "." } dependencies = [ { name = "applicationinsights" },