|
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,8 @@ 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 | + # Use an installed log handler if the caller set one; otherwise the interceptor's default. |
| 242 | + log_handler = get_log_handler() |
289 | 243 |
|
290 | 244 | self.logs_interceptor = UiPathRuntimeLogsInterceptor( |
291 | 245 | min_level=self.logs_min_level, |
@@ -359,6 +313,18 @@ def __exit__(self, exc_type, exc_val, exc_tb): |
359 | 313 | with open(self.output_file, "w") as f: |
360 | 314 | json.dump(output_payload, f, default=str) |
361 | 315 |
|
| 316 | + # Best-effort side channel: a sink failure must NOT reach the catch-all below, which would |
| 317 | + # rewrite the already-good output.json as FAULTED. |
| 318 | + result_sink = get_result_sink() |
| 319 | + if result_sink is not None and self.result.status in ( |
| 320 | + UiPathRuntimeStatus.SUCCESSFUL, |
| 321 | + UiPathRuntimeStatus.FAULTED, |
| 322 | + ): |
| 323 | + try: |
| 324 | + self._deliver_result(result_sink, output_payload) |
| 325 | + except Exception: |
| 326 | + logger.exception("Failed to deliver result to sink") |
| 327 | + |
362 | 328 | # Don't suppress exceptions |
363 | 329 | return False |
364 | 330 |
|
@@ -402,8 +368,16 @@ def __exit__(self, exc_type, exc_val, exc_tb): |
402 | 368 | # Restore original logging |
403 | 369 | if hasattr(self, "logs_interceptor"): |
404 | 370 | self.logs_interceptor.teardown() |
405 | | - if self.ipc_client is not None: |
406 | | - self.ipc_client.close() |
| 371 | + |
| 372 | + def _deliver_result(self, sink: "ResultSink", output_payload: Any) -> None: |
| 373 | + """Spill the output arguments to a file, then hand the result + that path to the sink.""" |
| 374 | + args_path = self.resolved_output_arguments_file_path |
| 375 | + # Avoid re-spilling if split_output_arguments already wrote this file. |
| 376 | + if not (self.split_output_arguments and self.job_id): |
| 377 | + os.makedirs(os.path.dirname(args_path), exist_ok=True) |
| 378 | + with open(args_path, "w") as f: |
| 379 | + json.dump(output_payload, f, default=str) |
| 380 | + sink(self.result, args_path) |
407 | 381 |
|
408 | 382 | @cached_property |
409 | 383 | def resolved_result_file_path(self) -> str: |
@@ -468,8 +442,6 @@ def with_defaults( |
468 | 442 | base.tenant_id = os.environ.get("UIPATH_TENANT_ID") |
469 | 443 | base.process_key = os.environ.get("UIPATH_PROCESS_UUID") |
470 | 444 | 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 | 445 |
|
474 | 446 | # Override with kwargs |
475 | 447 | for k, v in kwargs.items(): |
|
0 commit comments