Skip to content

Commit 0654cec

Browse files
robert-ursuclaude
andcommitted
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) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CFMKm7zHGcptS49z4cnnad
1 parent db4653c commit 0654cec

10 files changed

Lines changed: 793 additions & 63 deletions

File tree

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,99 @@
1+
"""Lets the server cancel the job that runs on its worker thread.
2+
3+
run/debug/eval drive their own event loop inside the ``asyncio.to_thread`` worker. That
4+
loop is the only place a cancellation can land, so the command publishes it here and
5+
the server reaches it through ``loop.call_soon_threadsafe``.
6+
7+
Kept import-light: ``_cli/__init__.py`` defers heavy imports.
8+
"""
9+
10+
import asyncio
11+
import contextvars
12+
import threading
13+
from typing import Any
14+
15+
CURRENT_JOB_CONTROL: contextvars.ContextVar["JobControl | None"] = (
16+
contextvars.ContextVar("uipath_current_job_control", default=None)
17+
)
18+
19+
20+
class JobControl:
21+
"""Handle on one job's event loop, shared by the server loop and the job thread."""
22+
23+
def __init__(self) -> None:
24+
"""Create an unbound control; the job binds its loop once it has one."""
25+
self.cancel_requested = False
26+
self._delivered = False
27+
self._loop: asyncio.AbstractEventLoop | None = None
28+
self._task: "asyncio.Task[Any] | None" = None
29+
self._sync = threading.Lock()
30+
31+
def bind(self, loop: asyncio.AbstractEventLoop, task: "asyncio.Task[Any]") -> None:
32+
"""Publish the job's loop and root task (job thread)."""
33+
with self._sync:
34+
self._loop, self._task = loop, task
35+
pending = self.cancel_requested
36+
if pending:
37+
self.cancel()
38+
39+
def unbind(self) -> None:
40+
"""Withdraw the loop once the root task is done (job thread)."""
41+
with self._sync:
42+
self._loop = self._task = None
43+
44+
def cancel(self) -> None:
45+
"""Cancel the job's root task, at most once.
46+
47+
A second delivery would land inside the runtime's cleanup ``finally`` blocks,
48+
the ones that write ``output.json``, and abort them. Before the job has a loop,
49+
the request is recorded and applied by ``bind``.
50+
"""
51+
with self._sync:
52+
self.cancel_requested = True
53+
if self._delivered or self._loop is None or self._task is None:
54+
return
55+
loop, task = self._loop, self._task
56+
self._delivered = True
57+
try:
58+
loop.call_soon_threadsafe(task.cancel)
59+
except RuntimeError:
60+
with self._sync:
61+
self._delivered = False
62+
63+
def cancel_all(self) -> None:
64+
"""Cancel every task on the job's loop, giving up on a clean cleanup."""
65+
with self._sync:
66+
self.cancel_requested = True
67+
loop = self._loop
68+
if loop is None:
69+
return
70+
71+
def _sweep() -> None:
72+
for task in asyncio.all_tasks(loop):
73+
task.cancel()
74+
75+
try:
76+
loop.call_soon_threadsafe(_sweep)
77+
except RuntimeError:
78+
pass
79+
80+
81+
def run_job_loop(coro: Any) -> Any:
82+
"""``asyncio.run`` that publishes its loop and root task to the job's control.
83+
84+
With no control in scope (``uipath run`` from a terminal) this is ``asyncio.run``.
85+
The loop is withdrawn before the runner closes, so a sweep can never cancel the
86+
runner's own wait on the job's executor threads.
87+
"""
88+
control = CURRENT_JOB_CONTROL.get()
89+
if control is None:
90+
return asyncio.run(coro)
91+
92+
with asyncio.Runner() as runner:
93+
loop = runner.get_loop()
94+
task = loop.create_task(coro, context=contextvars.copy_context())
95+
control.bind(loop, task)
96+
try:
97+
return loop.run_until_complete(task)
98+
finally:
99+
control.unbind()

0 commit comments

Comments
 (0)