|
3 | 3 | from types import TracebackType |
4 | 4 |
|
5 | 5 | from typing_extensions import Self |
| 6 | +from google.protobuf.message import Message |
6 | 7 |
|
7 | 8 | from a2a.client.middleware import ClientCallContext |
8 | 9 | from a2a.types.a2a_pb2 import ( |
@@ -158,3 +159,164 @@ async def get_extended_agent_card( |
158 | 159 | @abstractmethod |
159 | 160 | async def close(self) -> None: |
160 | 161 | """Closes the transport.""" |
| 162 | + |
| 163 | + |
| 164 | +class TenantTransportDecorator(ClientTransport): |
| 165 | + """A transport decorator that attaches a tenant to all requests.""" |
| 166 | + |
| 167 | + def __init__(self, base: ClientTransport, tenant: str): |
| 168 | + self._base = base |
| 169 | + self._tenant = tenant |
| 170 | + |
| 171 | + def update_tenant(self, request: Message) -> str | None: |
| 172 | + """Ensures the tenant is set on the request if provided and not already set. |
| 173 | +
|
| 174 | + Returns: |
| 175 | + The tenant used for the request. |
| 176 | + """ |
| 177 | + current_tenant = getattr(request, 'tenant', None) |
| 178 | + if current_tenant: |
| 179 | + return current_tenant |
| 180 | + |
| 181 | + if self._tenant and hasattr(request, 'tenant'): |
| 182 | + request.tenant = self._tenant |
| 183 | + return self._tenant |
| 184 | + return None |
| 185 | + |
| 186 | + async def send_message( |
| 187 | + self, |
| 188 | + request: SendMessageRequest, |
| 189 | + *, |
| 190 | + context: ClientCallContext | None = None, |
| 191 | + extensions: list[str] | None = None, |
| 192 | + ) -> SendMessageResponse: |
| 193 | + self.update_tenant(request) |
| 194 | + return await self._base.send_message( |
| 195 | + request, context=context, extensions=extensions |
| 196 | + ) |
| 197 | + |
| 198 | + async def send_message_streaming( |
| 199 | + self, |
| 200 | + request: SendMessageRequest, |
| 201 | + *, |
| 202 | + context: ClientCallContext | None = None, |
| 203 | + extensions: list[str] | None = None, |
| 204 | + ) -> AsyncGenerator[StreamResponse]: |
| 205 | + self.update_tenant(request) |
| 206 | + async for event in self._base.send_message_streaming( |
| 207 | + request, context=context, extensions=extensions |
| 208 | + ): |
| 209 | + yield event |
| 210 | + |
| 211 | + async def get_task( |
| 212 | + self, |
| 213 | + request: GetTaskRequest, |
| 214 | + *, |
| 215 | + context: ClientCallContext | None = None, |
| 216 | + extensions: list[str] | None = None, |
| 217 | + ) -> Task: |
| 218 | + self.update_tenant(request) |
| 219 | + return await self._base.get_task( |
| 220 | + request, context=context, extensions=extensions |
| 221 | + ) |
| 222 | + |
| 223 | + async def list_tasks( |
| 224 | + self, |
| 225 | + request: ListTasksRequest, |
| 226 | + *, |
| 227 | + context: ClientCallContext | None = None, |
| 228 | + extensions: list[str] | None = None, |
| 229 | + ) -> ListTasksResponse: |
| 230 | + self.update_tenant(request) |
| 231 | + return await self._base.list_tasks( |
| 232 | + request, context=context, extensions=extensions |
| 233 | + ) |
| 234 | + |
| 235 | + async def cancel_task( |
| 236 | + self, |
| 237 | + request: CancelTaskRequest, |
| 238 | + *, |
| 239 | + context: ClientCallContext | None = None, |
| 240 | + extensions: list[str] | None = None, |
| 241 | + ) -> Task: |
| 242 | + self.update_tenant(request) |
| 243 | + return await self._base.cancel_task( |
| 244 | + request, context=context, extensions=extensions |
| 245 | + ) |
| 246 | + |
| 247 | + async def set_task_callback( |
| 248 | + self, |
| 249 | + request: CreateTaskPushNotificationConfigRequest, |
| 250 | + *, |
| 251 | + context: ClientCallContext | None = None, |
| 252 | + extensions: list[str] | None = None, |
| 253 | + ) -> TaskPushNotificationConfig: |
| 254 | + self.update_tenant(request) |
| 255 | + return await self._base.set_task_callback( |
| 256 | + request, context=context, extensions=extensions |
| 257 | + ) |
| 258 | + |
| 259 | + async def get_task_callback( |
| 260 | + self, |
| 261 | + request: GetTaskPushNotificationConfigRequest, |
| 262 | + *, |
| 263 | + context: ClientCallContext | None = None, |
| 264 | + extensions: list[str] | None = None, |
| 265 | + ) -> TaskPushNotificationConfig: |
| 266 | + self.update_tenant(request) |
| 267 | + return await self._base.get_task_callback( |
| 268 | + request, context=context, extensions=extensions |
| 269 | + ) |
| 270 | + |
| 271 | + async def list_task_callback( |
| 272 | + self, |
| 273 | + request: ListTaskPushNotificationConfigsRequest, |
| 274 | + *, |
| 275 | + context: ClientCallContext | None = None, |
| 276 | + extensions: list[str] | None = None, |
| 277 | + ) -> ListTaskPushNotificationConfigsResponse: |
| 278 | + self.update_tenant(request) |
| 279 | + return await self._base.list_task_callback( |
| 280 | + request, context=context, extensions=extensions |
| 281 | + ) |
| 282 | + |
| 283 | + async def delete_task_callback( |
| 284 | + self, |
| 285 | + request: DeleteTaskPushNotificationConfigRequest, |
| 286 | + *, |
| 287 | + context: ClientCallContext | None = None, |
| 288 | + extensions: list[str] | None = None, |
| 289 | + ) -> None: |
| 290 | + self.update_tenant(request) |
| 291 | + await self._base.delete_task_callback( |
| 292 | + request, context=context, extensions=extensions |
| 293 | + ) |
| 294 | + |
| 295 | + async def subscribe( |
| 296 | + self, |
| 297 | + request: SubscribeToTaskRequest, |
| 298 | + *, |
| 299 | + context: ClientCallContext | None = None, |
| 300 | + extensions: list[str] | None = None, |
| 301 | + ) -> AsyncGenerator[StreamResponse]: |
| 302 | + self.update_tenant(request) |
| 303 | + async for event in self._base.subscribe( |
| 304 | + request, context=context, extensions=extensions |
| 305 | + ): |
| 306 | + yield event |
| 307 | + |
| 308 | + async def get_extended_agent_card( |
| 309 | + self, |
| 310 | + *, |
| 311 | + context: ClientCallContext | None = None, |
| 312 | + extensions: list[str] | None = None, |
| 313 | + signature_verifier: Callable[[AgentCard], None] | None = None, |
| 314 | + ) -> AgentCard: |
| 315 | + return await self._base.get_extended_agent_card( |
| 316 | + context=context, |
| 317 | + extensions=extensions, |
| 318 | + signature_verifier=signature_verifier, |
| 319 | + ) |
| 320 | + |
| 321 | + async def close(self) -> None: |
| 322 | + await self._base.close() |
0 commit comments