1- """The job-invocation IPC contract and the glue that routes a job's logs and result over it."""
1+ """The Python job-api IPC contract and the glue that routes a job's logs and result over it."""
22
33from __future__ import annotations
44
@@ -55,9 +55,11 @@ class ExecutorJobStatus(IntEnum):
5555
5656
5757@dataclass
58- class JobLogDto :
58+ class PythonJobLogDto :
5959 """A log entry; field names are the wire keys (do not rename)."""
6060
61+ JobKey : str
62+ ResumeVersion : int | None = None
6163 Message : str = ""
6264 LogLevel : int = LogLevel .INFORMATION .value
6365
@@ -74,26 +76,31 @@ class JobExecutorError:
7476
7577
7678@dataclass
77- class JobResultDto :
79+ class PythonJobResultDto :
7880 """The final result; field names are the wire keys (do not rename)."""
7981
80- id : str = ""
81- status : int = ExecutorJobStatus .SUCCESSFUL .value
82- outputArguments : Any = None
83- outputArgumentsFilePath : str | None = None
84- info : str | None = None
85- error : JobExecutorError | None = None
82+ JobKey : str
83+ ResumeVersion : int | None = None
84+ Status : int = ExecutorJobStatus .SUCCESSFUL .value
85+ OutputArguments : Any = None
86+ OutputArgumentsFilePath : str | None = None
87+ Info : str | None = None
88+ Error : JobExecutorError | None = None
8689
8790
88- class IJobInvocationCommonApi (ABC ):
89- """The job-invocation contract: logs + the final result. The class name is the endpoint key."""
91+ class IPythonJobApi (ABC ):
92+ """The Python job-api contract: logs + the final result. The class name is the endpoint key.
93+
94+ Every message names the run it belongs to (job key + resume version), so the peer can route a
95+ pooled callback to the right job and drop a straggler from a previous resume.
96+ """
9097
9198 @abstractmethod
92- async def SendLog (self , jobId : str , log : JobLogDto ) -> None :
99+ async def SendLog (self , log : PythonJobLogDto ) -> None :
93100 """Forward one log entry."""
94101
95102 @abstractmethod
96- async def SetResult (self , jobId : str , result : JobResultDto ) -> bool :
103+ async def SetResult (self , result : PythonJobResultDto ) -> bool :
97104 """Submit the final result."""
98105
99106
@@ -121,8 +128,11 @@ def _to_log_level(levelno: int) -> int:
121128
122129
123130def _to_result_dto (
124- job_id : str , result : Any , output_arguments_file_path : str
125- ) -> JobResultDto :
131+ job_key : str ,
132+ resume_version : int | None ,
133+ result : Any ,
134+ output_arguments_file_path : str ,
135+ ) -> PythonJobResultDto :
126136 error = None
127137 if result is not None and getattr (result , "error" , None ) is not None :
128138 category = result .error .category
@@ -135,11 +145,12 @@ def _to_result_dto(
135145 )
136146 raw_status = getattr (result , "status" , None )
137147 status_key = str (getattr (raw_status , "value" , raw_status ) or "successful" ).lower ()
138- return JobResultDto (
139- id = job_id ,
140- status = _EXECUTOR_STATUS .get (status_key , ExecutorJobStatus .SUCCESSFUL .value ),
141- outputArgumentsFilePath = output_arguments_file_path ,
142- error = error ,
148+ return PythonJobResultDto (
149+ JobKey = job_key ,
150+ ResumeVersion = resume_version ,
151+ Status = _EXECUTOR_STATUS .get (status_key , ExecutorJobStatus .SUCCESSFUL .value ),
152+ OutputArgumentsFilePath = output_arguments_file_path ,
153+ Error = error ,
143154 )
144155
145156
@@ -162,10 +173,15 @@ class _IpcLogHandler(logging.Handler):
162173 """Forwards each log record to the callback."""
163174
164175 def __init__ (
165- self , job_id : str , callback : Any , loop : asyncio .AbstractEventLoop
176+ self ,
177+ job_key : str ,
178+ resume_version : int | None ,
179+ callback : Any ,
180+ loop : asyncio .AbstractEventLoop ,
166181 ) -> None :
167182 super ().__init__ ()
168- self ._job_id = job_id
183+ self ._job_key = job_key
184+ self ._resume_version = resume_version
169185 self ._callback = callback
170186 self ._loop = loop
171187 self ._pending : set [Future [object ]] = set ()
@@ -177,11 +193,14 @@ def emit(self, record: logging.LogRecord) -> None:
177193 _to_original_stderr (self , record )
178194 return
179195 try :
180- dto = JobLogDto (
181- Message = self .format (record ), LogLevel = _to_log_level (record .levelno )
196+ dto = PythonJobLogDto (
197+ JobKey = self ._job_key ,
198+ ResumeVersion = self ._resume_version ,
199+ Message = self .format (record ),
200+ LogLevel = _to_log_level (record .levelno ),
182201 )
183202 future = asyncio .run_coroutine_threadsafe (
184- self ._callback .SendLog (self . _job_id , dto ), self ._loop
203+ self ._callback .SendLog (dto ), self ._loop
185204 )
186205 with self ._pending_lock :
187206 self ._pending .add (future )
@@ -226,9 +245,15 @@ async def aflush_pending(self, timeout: float = _LOG_FLUSH_TIMEOUT_S) -> None:
226245
227246
228247def install_runtime_sinks (
229- job_id : str , callback : Any , loop : asyncio .AbstractEventLoop
248+ job_key : str ,
249+ resume_version : int | None ,
250+ callback : Any ,
251+ loop : asyncio .AbstractEventLoop ,
230252) -> "_IpcLogHandler | None" :
231- """Install the log + result sinks, forwarding to ``callback`` on ``loop``.
253+ """Install the log + result sinks for the run ``(job_key, resume_version)``, forwarding to ``callback`` on ``loop``.
254+
255+ The peer routes a pooled callback by that pair exactly, so ``resume_version`` is the caller's
256+ decision: the value it was handed, or ``None`` on a lane that has none.
232257
233258 ``loop`` must run on a different thread than the one the sinks are invoked on, or the result ack
234259 deadlocks. Raises if this runtime has no sinks to install into.
@@ -246,15 +271,15 @@ def install_runtime_sinks(
246271 "Install uipath-runtime>=0.13.5."
247272 ) from e
248273
249- handler = _IpcLogHandler (job_id , callback , loop )
274+ handler = _IpcLogHandler (job_key , resume_version , callback , loop )
250275 handler .setFormatter (logging .Formatter ("%(message)s" ))
251276
252277 def _result_sink (result : Any , output_arguments_file_path : str ) -> None :
253- dto = _to_result_dto (job_id , result , output_arguments_file_path )
278+ dto = _to_result_dto (
279+ job_key , resume_version , result , output_arguments_file_path
280+ )
254281 try :
255- future = asyncio .run_coroutine_threadsafe (
256- callback .SetResult (job_id , dto ), loop
257- )
282+ future = asyncio .run_coroutine_threadsafe (callback .SetResult (dto ), loop )
258283 if future .result (timeout = _SET_RESULT_TIMEOUT_S ) is False :
259284 logger .error (
260285 "The handler rejected the job result (SetResult returned false)"
@@ -329,21 +354,23 @@ def _shutdown(self) -> None:
329354 _stop_loop_thread (self ._loop , self ._thread , _SET_RESULT_TIMEOUT_S )
330355
331356
332- def is_wire_job_id ( job_id : str | None ) -> TypeGuard [str ]:
333- """The peer types the job id as a Guid, and routes nothing for one it can't match."""
357+ def is_wire_job_key ( job_key : str | None ) -> TypeGuard [str ]:
358+ """The peer types the job key as a Guid, and routes nothing for one it can't match."""
334359 try :
335- parsed = uuid .UUID (str (job_id ))
360+ parsed = uuid .UUID (str (job_key ))
336361 except ValueError :
337362 return False
338363 # All-zeros is Guid's default, so it is the one unusable value a caller reaches by omission.
339364 return parsed .int != 0
340365
341366
342- def connect_handler_ipc (pipe : str , job_id : str | None ) -> _HandlerIpcConnection :
367+ def connect_handler_ipc (
368+ pipe : str , job_key : str | None , resume_version : int | None
369+ ) -> _HandlerIpcConnection :
343370 """Dial ``pipe`` on a dedicated loop/thread and install the sinks (see ``_HandlerIpcConnection``)."""
344- if not is_wire_job_id ( job_id ):
371+ if not is_wire_job_key ( job_key ):
345372 raise RuntimeError (
346- f"--handler-ipc-pipe needs UIPATH_JOB_KEY to be a job id ; got { job_id !r} ."
373+ f"--handler-ipc-pipe needs UIPATH_JOB_KEY to be a job key ; got { job_key !r} ."
347374 )
348375
349376 from uipath_ipc import IpcClient , NamedPipeClientTransport
@@ -369,7 +396,7 @@ async def _build() -> Any:
369396 request_timeout = _IPC_REQUEST_TIMEOUT_S ,
370397 max_message_size = _MAX_MESSAGE_BYTES ,
371398 )
372- proxy = client .get_proxy (IJobInvocationCommonApi ) # type: ignore[type-abstract]
399+ proxy = client .get_proxy (IPythonJobApi ) # type: ignore[type-abstract]
373400 return client , proxy
374401
375402 try :
@@ -380,7 +407,7 @@ async def _build() -> Any:
380407 _stop_loop_thread (loop , thread , _SET_RESULT_TIMEOUT_S )
381408 raise
382409 # Install in the caller's context: the sinks are contextvars, read in the job's own context at teardown.
383- handler = install_runtime_sinks (job_id , proxy , loop )
410+ handler = install_runtime_sinks (job_key , resume_version , proxy , loop )
384411 return _HandlerIpcConnection (client , loop , thread , handler )
385412
386413
@@ -392,10 +419,12 @@ async def disconnect_handler_ipc(conn: _HandlerIpcConnection) -> None:
392419
393420@contextlib .asynccontextmanager
394421async def handler_ipc_connection (
395- pipe : str | None , job_id : str | None
422+ pipe : str | None , job_key : str | None , resume_version : int | None
396423) -> AsyncIterator [Any ]:
397424 """Connect (if ``pipe`` is set) and always disconnect on exit; yields the connection or None."""
398- conn = connect_handler_ipc (pipe , job_id ) if pipe is not None else None
425+ conn = (
426+ connect_handler_ipc (pipe , job_key , resume_version ) if pipe is not None else None
427+ )
399428 try :
400429 yield conn
401430 finally :
0 commit comments