|
18 | 18 | UiPathErrorContract, |
19 | 19 | UiPathRuntimeError, |
20 | 20 | ) |
21 | | -from uipath.runtime.jobapi.client import IpcJobApiClient |
22 | | -from uipath.runtime.jobapi.log_handler import IpcSendLogHandler, PooledIpcSendLogHandler |
23 | | -from uipath.runtime.jobapi.pooled import get_pooled_log_sink |
24 | 21 | from uipath.runtime.logging._interceptor import UiPathRuntimeLogsInterceptor |
| 22 | +from uipath.runtime.output_sinks import ResultSink, get_log_handler, get_result_sink |
25 | 23 | from uipath.runtime.result import UiPathRuntimeResult, UiPathRuntimeStatus |
26 | 24 |
|
27 | 25 | logger = logging.getLogger(__name__) |
@@ -123,38 +121,8 @@ class UiPathRuntimeContext(BaseModel): |
123 | 121 | keep_state_file: bool = Field( |
124 | 122 | False, description="Prevents deletion of state file before running." |
125 | 123 | ) |
126 | | - ipc_endpoint: str | None = Field( |
127 | | - None, |
128 | | - description=( |
129 | | - "uipath-ipc endpoint (UIPATH_JOB_API_IPC_ENDPOINT) for streaming logs to " |
130 | | - "the handler in place of execution.log. The result stays on output.json." |
131 | | - ), |
132 | | - ) |
133 | | - ipc_job_id: str | None = Field( |
134 | | - None, |
135 | | - description="Handler job id (UIPATH_JOB_ID) used to route the IPC calls.", |
136 | | - ) |
137 | | - ipc_client: Any = Field(default=None, exclude=True, repr=False) |
138 | | - |
139 | 124 | model_config = ConfigDict(arbitrary_types_allowed=True, extra="allow") |
140 | 125 |
|
141 | | - @property |
142 | | - def ipc_active(self) -> bool: |
143 | | - """Whether logs flow over a per-job IPC pipe (non-pooled) instead of to files.""" |
144 | | - return bool(self.ipc_endpoint and self.ipc_job_id) |
145 | | - |
146 | | - @property |
147 | | - def pooled_ipc_active(self) -> bool: |
148 | | - """Whether logs stream over the pooled server's callback rather than a per-job pipe. |
149 | | -
|
150 | | - True when there is a job id but no endpoint, and the pooled server registered its sink. |
151 | | - """ |
152 | | - return ( |
153 | | - bool(self.ipc_job_id) |
154 | | - and not self.ipc_endpoint |
155 | | - and get_pooled_log_sink() is not None |
156 | | - ) |
157 | | - |
158 | 126 | def _apply_execution_source(self) -> None: |
159 | 127 | """Derive execution_source from the command, if not already set. |
160 | 128 |
|
@@ -270,22 +238,9 @@ def __enter__(self): |
270 | 238 | Returns: |
271 | 239 | The runtime context instance |
272 | 240 | """ |
273 | | - # Intercept all stdout/stderr/logs |
274 | | - # Write to file (runtime), stdout (debug) or log handler (if provided) |
275 | | - log_handler: logging.Handler | None = None |
276 | | - if self.ipc_active: |
277 | | - assert self.ipc_endpoint is not None and self.ipc_job_id is not None |
278 | | - self.ipc_client = IpcJobApiClient( |
279 | | - self.ipc_endpoint, self.ipc_job_id, logger |
280 | | - ) |
281 | | - self.ipc_client.start() |
282 | | - log_handler = IpcSendLogHandler(self.ipc_client) |
283 | | - log_handler.setFormatter(logging.Formatter("%(message)s")) |
284 | | - elif self.pooled_ipc_active: |
285 | | - sink = get_pooled_log_sink() |
286 | | - assert self.ipc_job_id is not None and sink is not None |
287 | | - log_handler = PooledIpcSendLogHandler(self.ipc_job_id, sink) |
288 | | - log_handler.setFormatter(logging.Formatter("%(message)s")) |
| 241 | + # Intercept all stdout/stderr/logs. A host may have installed a log handler (routing records |
| 242 | + # to IPC etc.); otherwise the interceptor opens execution.log. |
| 243 | + log_handler = get_log_handler() |
289 | 244 |
|
290 | 245 | self.logs_interceptor = UiPathRuntimeLogsInterceptor( |
291 | 246 | min_level=self.logs_min_level, |
@@ -359,6 +314,21 @@ def __exit__(self, exc_type, exc_val, exc_tb): |
359 | 314 | with open(self.output_file, "w") as f: |
360 | 315 | json.dump(output_payload, f, default=str) |
361 | 316 |
|
| 317 | + # Hand the result to a host-installed sink, if any. Successful/Faulted only — Suspended |
| 318 | + # keeps output.json (resume triggers live there); output.json is written above regardless. |
| 319 | + # Best-effort and ISOLATED: the sink is a side channel (e.g. IPC delivery), so a failure |
| 320 | + # here must not reach the catch-all below — that would rewrite the already-persisted good |
| 321 | + # output.json as FAULTED and fault a job that actually succeeded. |
| 322 | + result_sink = get_result_sink() |
| 323 | + if result_sink is not None and self.result.status in ( |
| 324 | + UiPathRuntimeStatus.SUCCESSFUL, |
| 325 | + UiPathRuntimeStatus.FAULTED, |
| 326 | + ): |
| 327 | + try: |
| 328 | + self._deliver_result(result_sink, output_payload) |
| 329 | + except Exception as sink_error: |
| 330 | + logger.error(f"Failed to deliver result to sink: {str(sink_error)}") |
| 331 | + |
362 | 332 | # Don't suppress exceptions |
363 | 333 | return False |
364 | 334 |
|
@@ -402,8 +372,21 @@ def __exit__(self, exc_type, exc_val, exc_tb): |
402 | 372 | # Restore original logging |
403 | 373 | if hasattr(self, "logs_interceptor"): |
404 | 374 | self.logs_interceptor.teardown() |
405 | | - if self.ipc_client is not None: |
406 | | - self.ipc_client.close() |
| 375 | + |
| 376 | + def _deliver_result(self, sink: "ResultSink", output_payload: Any) -> None: |
| 377 | + """Spill the output arguments to a file, then hand the result + that path to the host's sink. |
| 378 | +
|
| 379 | + The unbounded output arguments never travel inline — they go to a file the host delivers (it |
| 380 | + can stream it straight on, off-heap). The host maps the runtime result to its own wire shape. |
| 381 | + """ |
| 382 | + args_path = self.resolved_output_arguments_file_path |
| 383 | + # The split-output-arguments branch may already have spilled this exact payload to this exact |
| 384 | + # path; don't serialize the (potentially large) payload to disk a second time. |
| 385 | + if not (self.split_output_arguments and self.job_id): |
| 386 | + os.makedirs(os.path.dirname(args_path), exist_ok=True) |
| 387 | + with open(args_path, "w") as f: |
| 388 | + json.dump(output_payload, f, default=str) |
| 389 | + sink(self.result, args_path) |
407 | 390 |
|
408 | 391 | @cached_property |
409 | 392 | def resolved_result_file_path(self) -> str: |
@@ -468,8 +451,6 @@ def with_defaults( |
468 | 451 | base.tenant_id = os.environ.get("UIPATH_TENANT_ID") |
469 | 452 | base.process_key = os.environ.get("UIPATH_PROCESS_UUID") |
470 | 453 | base.folder_key = os.environ.get("UIPATH_FOLDER_KEY") |
471 | | - base.ipc_endpoint = os.environ.get("UIPATH_JOB_API_IPC_ENDPOINT") |
472 | | - base.ipc_job_id = os.environ.get("UIPATH_JOB_ID") |
473 | 454 |
|
474 | 455 | # Override with kwargs |
475 | 456 | for k, v in kwargs.items(): |
|
0 commit comments