Skip to content

Commit 3de1eda

Browse files
robert-ursuclaude
andcommitted
feat(cli): dispatch jobs asynchronously and push results/logs to the caller
`uipath server` used to block for a whole job and leave its outcome on disk for the caller to find. StartJob now enqueues the work and returns; the server pushes logs while the job runs and the terminal result when it finishes. Behaviour is a pure function of the request: a caller that supplies `resultCallbackSocket` gets async dispatch, one that does not gets the original blocking call, byte for byte. No handshake, no capability gate on the dispatch path -- an older caller cannot send the field and could not serve the callback if it did. POST {callback}/api/python/jobs/{jobKey}/result POST {callback}/api/python/jobs/{jobKey}/logs Execution stays serialised behind the process-wide lock: a job mutates process globals (logging handlers, OTel providers, env, cwd), so queueing changes who waits, not how many run. Real cancellation. run/debug/eval each drive their own event loop via asyncio.run inside the to_thread worker, so a job IS an event loop. They now go through `run_job_loop`, which publishes that loop to a JobControl carried on a ContextVar (to_thread propagates contextvars, so no monkeypatching). StopJob cancels the root task, which unwinds the runtime cooperatively -- its context managers still run, so output.json is still written and the file fallback still works. Escalates to a full loop sweep, then reports False rather than claiming a stop that did not happen. Also fixes two pre-existing defects that made the API result meaningless: _run_command_isolated hardcoded ExitCode 0 while click RETURNS ctx.exit(N)'s code under standalone_mode=False, so every ConsoleLogger.error path reported success; and the HTTP body now carries exitCode, which is the field the un-upgraded .NET handler already reads. Notes for review: - Server diagnostics go to a stderr handle bound at import, NOT ConsoleLogger. ConsoleLogger resolves sys.stdout at call time, and the runtime's interceptor has replaced it with a writer feeding the job's execution.log -- which the tailer then reads and posts back. With the callback down that is a self-feeding loop. - A 4xx from the callback is REJECTED, not UNREACHABLE: the caller is up and has moved on, so retrying cannot help and must not trip the shutdown path. - The log tailer holds back an unterminated tail; a handler writes the record and only then flushes, so a poll can otherwise split one line into two entries that cannot be rejoined. - StopJob takes resume_version as a trailing optional parameter rather than a DTO: uipath-ipc ignores a surplus wire arg and defaults a missing one, so old and new peers interoperate in both directions. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1 parent 6de5bed commit 3de1eda

14 files changed

Lines changed: 2486 additions & 35 deletions
Lines changed: 113 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,113 @@
1+
"""The seam that lets the server cancel the job running on its worker thread.
2+
3+
Every server command drives its own event loop with ``asyncio.run`` inside the
4+
``asyncio.to_thread`` worker. That loop is the only place a cancellation can actually
5+
land, so the command publishes it here and the server reaches back through
6+
``loop.call_soon_threadsafe``.
7+
8+
Deliberately import-light so ``_cli/__init__.py``'s lazy-import discipline is untouched.
9+
"""
10+
11+
import asyncio
12+
import contextvars
13+
import threading
14+
from typing import Any
15+
16+
CURRENT_JOB_CONTROL: contextvars.ContextVar["JobControl | None"] = (
17+
contextvars.ContextVar("uipath_current_job_control", default=None)
18+
)
19+
20+
21+
class JobControl:
22+
"""Handle on one job's event loop, shared between the server loop and the worker."""
23+
24+
def __init__(self, job_key: str) -> None:
25+
self.job_key = job_key
26+
self.cancel_requested = False
27+
self._loop: asyncio.AbstractEventLoop | None = None
28+
self._task: "asyncio.Task[Any] | None" = None
29+
# Bound from the job thread, read from the server loop.
30+
self._sync = threading.Lock()
31+
32+
# ---- job thread ---------------------------------------------------------
33+
34+
def bind(self, loop: asyncio.AbstractEventLoop, task: "asyncio.Task[Any]") -> None:
35+
with self._sync:
36+
self._loop, self._task = loop, task
37+
pending = self.cancel_requested
38+
if pending:
39+
# A stop that raced the loop's creation: apply it now rather than losing it.
40+
self.cancel()
41+
42+
def unbind(self) -> None:
43+
with self._sync:
44+
self._loop = self._task = None
45+
46+
# ---- server loop --------------------------------------------------------
47+
48+
@property
49+
def bound(self) -> bool:
50+
with self._sync:
51+
return self._task is not None
52+
53+
def cancel(self) -> bool:
54+
"""Deliver CancelledError to the job's ROOT task only.
55+
56+
Call this at most once. A second cancel lands inside the runtime's cleanup
57+
``finally`` blocks — the ones that write ``output.json`` and tear the log
58+
interceptor down — and aborts them, which would destroy the very fallback the
59+
caller relies on. Escalate with ``cancel_all`` instead.
60+
"""
61+
with self._sync:
62+
self.cancel_requested = True
63+
loop, task = self._loop, self._task
64+
if loop is None or task is None:
65+
# Not bound yet; bind() will apply it.
66+
return False
67+
try:
68+
loop.call_soon_threadsafe(task.cancel)
69+
except RuntimeError:
70+
# Loop already closed — the job is finishing anyway.
71+
return False
72+
return True
73+
74+
def cancel_all(self) -> bool:
75+
"""Escalation: cancel every task on the job's loop.
76+
77+
This includes whatever the cleanup is awaiting, so it gives up on a clean
78+
``output.json``. Only worth doing once a polite cancel has already failed.
79+
"""
80+
with self._sync:
81+
loop = self._loop
82+
if loop is None:
83+
return False
84+
85+
def _sweep() -> None:
86+
for task in asyncio.all_tasks(loop):
87+
task.cancel()
88+
89+
try:
90+
loop.call_soon_threadsafe(_sweep)
91+
except RuntimeError:
92+
return False
93+
return True
94+
95+
96+
def run_job_loop(coro: Any) -> Any:
97+
"""``asyncio.run`` that publishes its loop and root task to the JobControl in scope.
98+
99+
Outside the server — ``uipath run`` on a terminal — there is no handle in scope and
100+
this is ``asyncio.run`` verbatim.
101+
"""
102+
control = CURRENT_JOB_CONTROL.get()
103+
if control is None:
104+
return asyncio.run(coro)
105+
106+
with asyncio.Runner() as runner:
107+
loop = runner.get_loop()
108+
task = loop.create_task(coro, context=contextvars.copy_context())
109+
control.bind(loop, task)
110+
try:
111+
return loop.run_until_complete(task)
112+
finally:
113+
control.unbind()

‎packages/uipath/src/uipath/_cli/_server_core.py‎

Lines changed: 168 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,13 @@
11
"""Transport-agnostic job core shared by the HTTP and uipath-ipc channels."""
22

33
import asyncio
4+
import json
45
import os
56
import shlex
7+
from collections.abc import Callable
68
from typing import Any
79

10+
from ._job_control import CURRENT_JOB_CONTROL, JobControl
811
from .cli_debug import debug
912
from .cli_eval import eval
1013
from .cli_run import run
@@ -45,17 +48,174 @@ def parse_args(args: str | list[str] | None) -> list[str]:
4548
return []
4649

4750

51+
# The document is carried inline on the result push when it fits. uipath_ipc caps a
52+
# frame at 2 MiB and Orchestrator already spills large outputs to an attachment, so
53+
# anything bigger stays on disk and the caller reads the file as it always has.
54+
MAX_INLINE_DOCUMENT_BYTES = 1024 * 1024
55+
56+
DEFAULT_RUNTIME_DIR = "__uipath"
57+
DEFAULT_RESULT_FILE = "output.json"
58+
DEFAULT_LOGS_FILE = "execution.log"
59+
60+
# Distinct from any click exit code, so the caller can tell "stopped on request" from
61+
# "the job failed" and report Stopped rather than Faulted.
62+
EXIT_CODE_STOPPED = 143
63+
64+
65+
def _resolve_runtime_file(
66+
config_path: str, base_dir: str, key: str, default_name: str
67+
) -> str | None:
68+
"""Resolve one ``runtime.*`` file path from a uipath.json.
69+
70+
Mirrors UiPathRuntimeContext.from_config's ``runtime.dir`` / ``runtime.<key>``
71+
mapping. ``base_dir`` anchors relative paths so this works without ever changing
72+
the process's cwd.
73+
"""
74+
if not os.path.isabs(config_path):
75+
config_path = os.path.join(base_dir, config_path)
76+
77+
runtime: dict[str, Any] = {}
78+
try:
79+
with open(config_path, encoding="utf-8") as f:
80+
loaded = json.load(f)
81+
if isinstance(loaded, dict):
82+
runtime = loaded.get("runtime") or {}
83+
except (OSError, json.JSONDecodeError, UnicodeDecodeError):
84+
# Fall through to the runtime's own defaults rather than giving up. Returning
85+
# None here would stop the log tailer from starting at all — and by then the
86+
# caller has already been told we forward logs and stopped tailing the file
87+
# itself, so the job would produce no logs anywhere.
88+
runtime = {}
89+
90+
directory = runtime.get("dir") or DEFAULT_RUNTIME_DIR
91+
name = runtime.get(key) or default_name
92+
if not isinstance(directory, str) or not isinstance(name, str):
93+
directory, name = DEFAULT_RUNTIME_DIR, default_name
94+
95+
if not os.path.isabs(directory):
96+
directory = os.path.join(base_dir, directory)
97+
return os.path.abspath(os.path.join(directory, name))
98+
99+
100+
def resolve_result_file_path() -> str | None:
101+
"""Where this job's terminal document lives. Call inside the job's env and cwd."""
102+
return _resolve_runtime_file(
103+
os.environ.get("UIPATH_CONFIG_PATH", "uipath.json"),
104+
os.getcwd(),
105+
"outputFile",
106+
DEFAULT_RESULT_FILE,
107+
)
108+
109+
110+
def resolve_logs_file_path(
111+
env_vars: dict[str, str] | None, working_dir: str | None
112+
) -> str | None:
113+
"""Where this job will write its log file, derived from the request alone.
114+
115+
Pure with respect to process state: the log tailer has to know the path *before*
116+
the job takes the lock and applies its env/cwd.
117+
"""
118+
env_vars = env_vars or {}
119+
base_dir = working_dir or os.getcwd()
120+
config_path = env_vars.get("UIPATH_CONFIG_PATH", "uipath.json")
121+
return _resolve_runtime_file(config_path, base_dir, "logsFile", DEFAULT_LOGS_FILE)
122+
123+
124+
def _read_result_document() -> tuple[str | None, str]:
125+
"""Return ``(document, conveyance)`` for the terminal result document.
126+
127+
``conveyance`` is ``inline`` when the document rides the wire, ``file`` when the
128+
caller must read it from disk (too large, unreadable, or never written).
129+
"""
130+
path = resolve_result_file_path()
131+
if not path or not os.path.exists(path):
132+
return None, "file"
133+
134+
try:
135+
if os.path.getsize(path) > MAX_INLINE_DOCUMENT_BYTES:
136+
return None, "file"
137+
with open(path, encoding="utf-8") as f:
138+
return f.read(), "inline"
139+
except (OSError, UnicodeDecodeError):
140+
return None, "file"
141+
142+
143+
async def _invoke_command(
144+
cmd: Any, args: list[str], control: "JobControl | None" = None
145+
) -> dict[str, Any]:
146+
"""Invoke one click command and classify how it ended."""
147+
# asyncio.to_thread propagates contextvars, so the command reads this inside the
148+
# worker and publishes its event loop back through it.
149+
token = CURRENT_JOB_CONTROL.set(control) if control is not None else None
150+
try:
151+
result_value = await asyncio.to_thread(cmd.main, args, standalone_mode=False)
152+
# Under standalone_mode=False click RETURNS ctx.exit(N)'s code instead of
153+
# raising SystemExit, so a bare int is the exit code, not a result — every
154+
# ConsoleLogger.error path lands here via ctx.exit(1). The run/debug/eval
155+
# callbacks only ever return a result object or None, so this is unambiguous.
156+
if isinstance(result_value, int) and not isinstance(result_value, bool):
157+
return {
158+
"ExitCode": result_value,
159+
"Error": None if result_value == 0 else f"Exit code: {result_value}",
160+
"Result": None,
161+
"Unexpected": False,
162+
}
163+
return {
164+
"ExitCode": 0,
165+
"Error": None,
166+
"Result": result_value,
167+
"Unexpected": False,
168+
}
169+
except SystemExit as e:
170+
exit_code = e.code if isinstance(e.code, int) else 1
171+
return {
172+
"ExitCode": exit_code,
173+
"Error": None if exit_code == 0 else f"Exit code: {exit_code}",
174+
"Result": None,
175+
"Unexpected": False,
176+
}
177+
except asyncio.CancelledError:
178+
# Two very different things arrive here. cancelling() > 0 means OUR awaiting task
179+
# was cancelled (server shutdown) — never swallow that, or the shutdown stalls.
180+
# 0 means the CancelledError travelled out of the job's own loop, i.e. our StopJob
181+
# landed: that is an OUTCOME, and swallowing it keeps this task alive so the lock
182+
# unwinds and the result document still gets read.
183+
current = asyncio.current_task()
184+
if current is not None and current.cancelling() > 0:
185+
raise
186+
return {
187+
"ExitCode": EXIT_CODE_STOPPED,
188+
"Error": "Job stopped on request",
189+
"Result": None,
190+
"Unexpected": False,
191+
"Stopped": True,
192+
}
193+
except Exception as e: # report any job failure as a result, not a fault
194+
return {"ExitCode": 1, "Error": str(e), "Result": None, "Unexpected": True}
195+
finally:
196+
if token is not None:
197+
CURRENT_JOB_CONTROL.reset(token)
198+
199+
48200
async def _run_command_isolated(
49201
cmd: Any,
50202
args: list[str],
51203
env_vars: dict[str, str],
52204
working_dir: str | None,
205+
on_started: Callable[[], None] | None = None,
206+
control: "JobControl | None" = None,
53207
) -> dict[str, Any]:
54-
"""Run one command with per-job env/cwd isolation (the shared job core)."""
208+
"""Run one command with per-job env/cwd isolation (the shared job core).
209+
210+
``on_started`` fires once the lock is held, i.e. the moment the job stops being
211+
queued and becomes uncancellable.
212+
"""
55213
if _state.lock is None or _state.baseline_env is None:
56214
raise RuntimeError("Server state not initialized")
57215

58216
async with _state.lock:
217+
if on_started is not None:
218+
on_started()
59219
original_cwd = os.getcwd()
60220
try:
61221
# Start from server baseline + request env vars only, so nothing from
@@ -79,25 +239,13 @@ async def _run_command_isolated(
79239
"ClientError": True,
80240
}
81241

82-
result_value = await asyncio.to_thread(
83-
cmd.main, args, standalone_mode=False
84-
)
85-
return {
86-
"ExitCode": 0,
87-
"Error": None,
88-
"Result": result_value,
89-
"Unexpected": False,
90-
}
91-
except SystemExit as e:
92-
exit_code = e.code if isinstance(e.code, int) else 1
93-
return {
94-
"ExitCode": exit_code,
95-
"Error": None if exit_code == 0 else f"Exit code: {exit_code}",
96-
"Result": None,
97-
"Unexpected": False,
98-
}
99-
except Exception as e: # report any job failure as a result, not a fault
100-
return {"ExitCode": 1, "Error": str(e), "Result": None, "Unexpected": True}
242+
outcome = await _invoke_command(cmd, args, control)
243+
# Must happen before the finally below restores env/cwd: the document's
244+
# location comes from this job's UIPATH_CONFIG_PATH and may be relative.
245+
document, conveyance = _read_result_document()
246+
outcome["Document"] = document
247+
outcome["DocumentConveyance"] = conveyance
248+
return outcome
101249
finally:
102250
# Restore to server baseline.
103251
try:

0 commit comments

Comments
 (0)