diff --git a/packages/uipath/pyproject.toml b/packages/uipath/pyproject.toml index bb12bfa6f..1ce573930 100644 --- a/packages/uipath/pyproject.toml +++ b/packages/uipath/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "uipath" -version = "2.14.25" +version = "2.14.26" 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_api.py b/packages/uipath/src/uipath/_cli/_job_api.py index 8777bc3a8..223b41894 100644 --- a/packages/uipath/src/uipath/_cli/_job_api.py +++ b/packages/uipath/src/uipath/_cli/_job_api.py @@ -120,6 +120,7 @@ def _to_log_level(levelno: int) -> int: # Defensive only: the runtime delivers a result for SUCCESSFUL and FAULTED alone, and the peer # resolves a suspended job from output.json (resume triggers never cross the wire). "suspended": ExecutorJobStatus.SUSPENDED.value, + "stopped": ExecutorJobStatus.STOPPED.value, } @@ -144,7 +145,8 @@ def _to_result_dto( return PythonJobResultDto( jobKey=job_key, resumeVersion=resume_version, - status=_EXECUTOR_STATUS.get(status_key, ExecutorJobStatus.SUCCESSFUL.value), + # A status this build does not know must not be reported as a success. + status=_EXECUTOR_STATUS.get(status_key, ExecutorJobStatus.FAULTED.value), outputArgumentsFilePath=output_arguments_file_path, error=error, ) 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 2ff91b6f6..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]: @@ -112,6 +156,163 @@ def parse_args(args: str | list[str] | None) -> list[str]: return [] +def _exit_code_outcome(exit_code: int) -> dict[str, Any]: + return { + "ExitCode": exit_code, + "Error": None if exit_code == 0 else f"Exit code: {exit_code}", + "Result": None, + "Unexpected": False, + } + + +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], @@ -119,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 @@ -160,26 +394,12 @@ 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() - return { - "ExitCode": 0, - "Error": None, - "Result": result_value, - "Unexpected": False, - } except SystemExit as e: - exit_code = e.code if isinstance(e.code, int) else 1 - return { - "ExitCode": exit_code, - "Error": None if exit_code == 0 else f"Exit code: {exit_code}", - "Result": None, - "Unexpected": False, - } + 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 return {"ExitCode": 1, "Error": str(e), "Result": None, "Unexpected": True} finally: @@ -190,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 66bdfad10..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())) @@ -538,6 +539,7 @@ async def execute_eval(): "uipath.json spec:", "https://github.com/UiPath/uipath-python/blob/main/packages/uipath/specs/uipath.spec.md", ) + click.get_current_context().exit(1) except ValueError as e: console.error(str(e)) except Exception as e: diff --git a/packages/uipath/src/uipath/_cli/cli_run.py b/packages/uipath/src/uipath/_cli/cli_run.py index d98b92653..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())) @@ -377,7 +377,7 @@ async def execute() -> None: "uipath.json spec:", "https://github.com/UiPath/uipath-python/blob/main/packages/uipath/specs/uipath.spec.md", ) - return + click.get_current_context().exit(1) except UiPathRuntimeError as e: console.error(f"{e.error_info.title} - {e.error_info.detail}") except Exception as e: diff --git a/packages/uipath/src/uipath/_cli/cli_server.py b/packages/uipath/src/uipath/_cli/cli_server.py index dc5e31bae..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,30 +185,106 @@ 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"] if result["Unexpected"]: return web.json_response( - {"success": False, "job_key": job_key, "error": result["Error"]}, + { + "success": False, + "job_key": job_key, + "exitCode": exit_code, + "error": result["Error"], + }, status=500, ) if result.get("ClientError"): # Request-shaped failure (e.g. bad working directory) — 4xx, not 200. return web.json_response( - {"success": False, "job_key": job_key, "error": result["Error"]}, + { + "success": False, + "job_key": job_key, + "exitCode": exit_code, + "error": result["Error"], + }, status=400, ) - if result["ExitCode"] == 0: + if exit_code == 0: return web.json_response( - {"success": True, "job_key": job_key, "result": result["Result"]} + { + "success": True, + "job_key": job_key, + "exitCode": exit_code, + "result": result["Result"], + } ) return web.json_response( - {"success": False, "job_key": job_key, "error": result["Error"]} + { + "success": False, + "job_key": job_key, + "exitCode": exit_code, + "error": result["Error"], + } ) +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]"} @@ -242,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 586680274..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. @@ -162,9 +170,11 @@ def _install() -> None: async def StopJob(self, request: PythonServerStopJobRequest) -> bool: console.info( f"StopJob requested for {_run_id(request.jobKey, request.resumeVersion)} " - f"(force={request.forceStop}) (no-op)" + f"(force={request.forceStop})" + ) + return await stop_job( + request.jobKey, request.resumeVersion, force=request.forceStop ) - return True async def start_ipc_server(pipe_name: str) -> None: diff --git a/packages/uipath/tests/cli/eval/test_eval_discovery.py b/packages/uipath/tests/cli/eval/test_eval_discovery.py index bce44f86c..743fc848e 100644 --- a/packages/uipath/tests/cli/eval/test_eval_discovery.py +++ b/packages/uipath/tests/cli/eval/test_eval_discovery.py @@ -83,7 +83,7 @@ def test_multiple_entrypoints_shows_usage_help( ): result = runner.invoke(cli, ["eval"]) - assert result.exit_code == 0 + assert result.exit_code == 1 assert "Available entrypoints:" in result.output assert "agent_a" in result.output assert "agent_b" in result.output @@ -116,7 +116,7 @@ def test_multiple_entrypoints_no_eval_sets(self, runner: CliRunner, temp_dir: st ): result = runner.invoke(cli, ["eval"]) - assert result.exit_code == 0 + assert result.exit_code == 1 assert "Available entrypoints:" in result.output assert "a" in result.output assert "b" in result.output @@ -165,7 +165,7 @@ def test_multiple_eval_sets_shows_usage_help( ): result = runner.invoke(cli, ["eval"]) - assert result.exit_code == 0 + assert result.exit_code == 1 assert "Available entrypoints:" in result.output assert "my_agent" in result.output assert "Available eval sets:" in result.output @@ -302,7 +302,7 @@ def test_no_entrypoints_shows_helpful_message( ): result = runner.invoke(cli, ["eval"]) - assert result.exit_code == 0 + assert result.exit_code == 1 assert "No entrypoints found" in result.output assert "Usage: uipath eval " in result.output @@ -349,7 +349,7 @@ def test_explicit_entrypoint_skips_entrypoint_discovery( result = runner.invoke(cli, ["eval", "agent_a"]) # Should still show usage help because multiple eval sets - assert result.exit_code == 0 + assert result.exit_code == 1 assert "Available eval sets:" in result.output assert "set-a.json" in result.output assert "set-b.json" in result.output diff --git a/packages/uipath/tests/cli/test_job_api.py b/packages/uipath/tests/cli/test_job_api.py index df2608890..913ae046e 100644 --- a/packages/uipath/tests/cli/test_job_api.py +++ b/packages/uipath/tests/cli/test_job_api.py @@ -82,6 +82,24 @@ class _Result: assert dto.status == _job_api.ExecutorJobStatus.SUSPENDED.value +def test_to_result_dto_maps_stopped(): + class _Result: + status = "stopped" + error = None + + dto = _job_api._to_result_dto("j", None, _Result(), "p.args") + assert dto.status == _job_api.ExecutorJobStatus.STOPPED.value + + +def test_to_result_dto_reports_an_unknown_status_as_faulted(): + class _Result: + status = "something-new" + error = None + + dto = _job_api._to_result_dto("j", None, _Result(), "p.args") + assert dto.status == _job_api.ExecutorJobStatus.FAULTED.value + + def test_to_log_level_maps_python_levels_to_wire_values(): assert _job_api._to_log_level(logging.CRITICAL) == _job_api.LogLevel.CRITICAL assert _job_api._to_log_level(logging.ERROR) == _job_api.LogLevel.ERROR diff --git a/packages/uipath/tests/cli/test_run.py b/packages/uipath/tests/cli/test_run.py index a47e5a4d7..4370cd42f 100644 --- a/packages/uipath/tests/cli/test_run.py +++ b/packages/uipath/tests/cli/test_run.py @@ -359,7 +359,7 @@ def test_no_entrypoint_multiple_available( ): result = runner.invoke(cli, ["run"]) - assert result.exit_code == 0 + assert result.exit_code == 1 assert "Available entrypoints:" in result.output assert "agent_a" in result.output assert "agent_b" in result.output @@ -383,7 +383,7 @@ def test_no_entrypoint_none_available(self, runner: CliRunner, temp_dir: str): ): result = runner.invoke(cli, ["run"]) - assert result.exit_code == 0 + assert result.exit_code == 1 assert "No entrypoints found" in result.output assert "Usage: uipath run" in result.output mock_factory.new_runtime.assert_not_awaited() diff --git a/packages/uipath/tests/cli/test_server.py b/packages/uipath/tests/cli/test_server.py index 70d29cb38..0de506e2f 100644 --- a/packages/uipath/tests/cli/test_server.py +++ b/packages/uipath/tests/cli/test_server.py @@ -132,12 +132,34 @@ def test_start_job_success(self, server, temp_dir, simple_script): assert response["success"] is True assert response["job_key"] == job_key + assert response["exitCode"] == 0 assert os.path.exists(output_file) with open(output_file, "r") as f: output = f.read() assert "Hello" in output + def test_failing_job_reports_its_exit_code(self, server, temp_dir): + """A job that fails through ConsoleLogger.error must not be reported as a success.""" + port = server + + with pytest.MonkeyPatch().context() as mp: + mp.chdir(temp_dir) + + script_file = "entrypoint.py" + with open(os.path.join(temp_dir, script_file), "w") as f: + f.write("def main(input: dict) -> str:\n raise ValueError('boom')\n") + + with open(os.path.join(temp_dir, "uipath.json"), "w") as f: + json.dump(create_uipath_json(script_file), f) + + response = asyncio.run( + start_job(port, "failing-job", "run", ["main", "{}"]) + ) + + assert response["success"] is False + assert response["exitCode"] == 1 + def test_start_job_unknown_command(self, server): """Test starting a job with unknown command.""" port = server 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 2d7d1b8ae..18798d44a 100644 --- a/packages/uipath/tests/cli/test_server_ipc.py +++ b/packages/uipath/tests/cli/test_server_ipc.py @@ -191,8 +191,7 @@ def test_stop_job_accepts_resume_version_and_force_stop(self, pipe): assert result is True - def test_stop_job_returns_true(self, pipe): - """StopJob is a no-op stub today, but must ack (bool) so the call is awaitable.""" + 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}) @@ -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 7f918fe1d..67e764613 100644 --- a/packages/uipath/tests/cli/test_server_job_core.py +++ b/packages/uipath/tests/cli/test_server_job_core.py @@ -12,9 +12,11 @@ from typing import Any from unittest.mock import Mock +import click import pytest from uipath._cli import _server_core +from uipath._cli._utils._console import ConsoleLogger @pytest.fixture @@ -61,6 +63,37 @@ async def test_maps_system_exit_code(restore_state: Any) -> None: assert result["Unexpected"] is False +async def test_maps_a_returned_click_exit_code(restore_state: Any) -> None: + _init(restore_state) + + @click.command() + def failing() -> None: + ConsoleLogger().error("boom") + + result = await _server_core._run_command_isolated(failing, [], {}, None) + assert result["ExitCode"] == 1 + assert result["Error"] == "Exit code: 1" + assert result["Unexpected"] is False + + +async def test_a_returned_zero_exit_code_is_success(restore_state: Any) -> None: + _init(restore_state) + cmd = Mock() + cmd.main.return_value = 0 + result = await _server_core._run_command_isolated(cmd, [], {}, None) + assert result["ExitCode"] == 0 + assert result["Error"] is None + + +async def test_a_non_int_return_value_is_the_result(restore_state: Any) -> None: + _init(restore_state) + cmd = Mock() + cmd.main.return_value = True + result = await _server_core._run_command_isolated(cmd, [], {}, None) + assert result["ExitCode"] == 0 + assert result["Result"] is True + + async def test_reports_unexpected_exception(restore_state: Any) -> None: _init(restore_state) cmd = Mock() @@ -427,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 bc92233bb..ebf95dbc7 100644 --- a/packages/uipath/uv.lock +++ b/packages/uipath/uv.lock @@ -2599,7 +2599,7 @@ wheels = [ [[package]] name = "uipath" -version = "2.14.25" +version = "2.14.26" source = { editable = "." } dependencies = [ { name = "applicationinsights" },