|
| 1 | +import json |
1 | 2 | import logging |
| 3 | +from typing import Any, Dict |
| 4 | + |
| 5 | +import httpx |
2 | 6 |
|
3 | 7 | from .._config import Config |
4 | 8 | from .._execution_context import ExecutionContext |
5 | 9 | from .._utils import Endpoint, RequestSpec |
6 | | -from ..models import Connection, ConnectionToken |
| 10 | +from .._utils._ssl_context import get_httpx_client_kwargs |
| 11 | +from ..models import Connection, ConnectionToken, EventArguments |
7 | 12 | from ..tracing._traced import traced |
8 | 13 | from ._base_service import BaseService |
9 | 14 |
|
@@ -112,6 +117,84 @@ async def retrieve_token_async(self, key: str) -> ConnectionToken: |
112 | 117 | ) |
113 | 118 | return ConnectionToken.model_validate(response.json()) |
114 | 119 |
|
| 120 | + @traced( |
| 121 | + name="connections_retrieve_event_payload", |
| 122 | + run_type="uipath", |
| 123 | + ) |
| 124 | + def retrieve_event_payload(self, event_args: EventArguments) -> Dict[str, Any]: |
| 125 | + """Retrieve event payload from UiPath Integration Service. |
| 126 | +
|
| 127 | + Args: |
| 128 | + key (str): The unique identifier of the connection. |
| 129 | + event_args (EventArguments): The event arguments. Should be passed along from the job's input. |
| 130 | +
|
| 131 | + Returns: |
| 132 | + Dict[str, Any]: The event payload data |
| 133 | + """ |
| 134 | + if not event_args.additional_event_data: |
| 135 | + raise ValueError("additional_event_data is required") |
| 136 | + |
| 137 | + # Parse additional event data to get event id |
| 138 | + event_data = json.loads(event_args.additional_event_data) |
| 139 | + |
| 140 | + event_id = None |
| 141 | + if "processedEventId" in event_data: |
| 142 | + event_id = event_data["processedEventId"] |
| 143 | + elif "rawEventId" in event_data: |
| 144 | + event_id = event_data["rawEventId"] |
| 145 | + else: |
| 146 | + raise ValueError("Event Id not found in additional event data") |
| 147 | + |
| 148 | + # Build request URL using connection token's API base URI |
| 149 | + spec = self._retrieve_event_payload_spec("v1", event_id) |
| 150 | + |
| 151 | + response = self.request(spec.method, url=spec.endpoint) |
| 152 | + |
| 153 | + return response.json() |
| 154 | + |
| 155 | + @traced( |
| 156 | + name="connections_retrieve_event_payload", |
| 157 | + run_type="uipath", |
| 158 | + ) |
| 159 | + async def retrieve_event_payload_async( |
| 160 | + self, event_args: EventArguments |
| 161 | + ) -> Dict[str, Any]: |
| 162 | + """Retrieve event payload from UiPath Integration Service. |
| 163 | +
|
| 164 | + Args: |
| 165 | + key (str): The unique identifier of the connection. |
| 166 | + event_args (EventArguments): The event arguments. Should be passed along from the job's input. |
| 167 | +
|
| 168 | + Returns: |
| 169 | + Dict[str, Any]: The event payload data |
| 170 | + """ |
| 171 | + if not event_args.additional_event_data: |
| 172 | + raise ValueError("additional_event_data is required") |
| 173 | + |
| 174 | + # Parse additional event data to get event id |
| 175 | + event_data = json.loads(event_args.additional_event_data) |
| 176 | + |
| 177 | + event_id = None |
| 178 | + if "processedEventId" in event_data: |
| 179 | + event_id = event_data["processedEventId"] |
| 180 | + elif "rawEventId" in event_data: |
| 181 | + event_id = event_data["rawEventId"] |
| 182 | + else: |
| 183 | + raise ValueError("Event Id not found in additional event data") |
| 184 | + |
| 185 | + # Build request URL using connection token's API base URI |
| 186 | + spec = self._retrieve_event_payload_spec("v1", event_id) |
| 187 | + |
| 188 | + response = await self.request_async(spec.method, url=spec.endpoint) |
| 189 | + |
| 190 | + return response.json() |
| 191 | + |
| 192 | + def _retrieve_event_payload_spec(self, version: str, event_id: str) -> RequestSpec: |
| 193 | + return RequestSpec( |
| 194 | + method="GET", |
| 195 | + endpoint=Endpoint(f"/elements_/{version}/events/{event_id}"), |
| 196 | + ) |
| 197 | + |
115 | 198 | def _retrieve_spec(self, key: str) -> RequestSpec: |
116 | 199 | return RequestSpec( |
117 | 200 | method="GET", |
|
0 commit comments