airbyte.cloud.client
PyAirbyte Cloud client.
Organization and workspace resolution
Most operations need an organization and/or a workspace context. CloudClient derives
that context from what the caller supplies, falling back to the credentials' own scope,
the authenticated user's default workspace, and finally to the authenticated user's
organization memberships.
Two rules govern the whole protocol: an explicitly passed ID always beats an ambient one, and within each of those tiers an organization ID beats a workspace ID.
Where a workspace ID comes from
A workspace ID reaches the client from one of four places, in order:
- The
workspace_idargument on the operation, such asCloudClient.get_workspace. - The
X-Airbyte-Workspace-Idheader, when running as an MCP server over HTTP. - The
AIRBYTE_CLOUD_WORKSPACE_ID(orAIRBYTE_WORKSPACE_ID) environment variable, read when the client is built withCloudClient.from_auth(env_vars=True). - The authenticated user's default workspace from their Airbyte user record.
The configured workspace becomes CloudClient.default_workspace_id, the ambient
workspace context for the client. Workspace-scoped operations use it, then the
authenticated user's default workspace, whenever no workspace_id argument is passed.
CloudClient.get_workspace raises when neither is available.
If a workspace ID is known
The workspace determines the organization: its parent organization is fetched in a single call and used as the organization context. That covers the common case, and nothing below applies.
The one exception is an explicit organization_id or organization_name argument,
which always wins over a workspace-derived organization — as does a configured
CloudClient.organization_id over an ambient workspace.
If an organization ID is known but no workspace ID
The organization is used as-is, whether it came from the organization_id argument,
from organization_name (an exact-name lookup, so it is never used to infer a
default), or from the credentials as CloudClient.organization_id.
If neither is known
CloudClient.list_workspaces uses privilege_scope to choose the search breadth:
MEMBER_OF lists direct workspace grants, ORGANIZATION_ADMIN lists workspaces in
member organizations, INSTANCE_ADMIN lists every workspace for instance admins, and
ANY chooses the broadest scope available to the caller. CloudClient.get_organization
called with no arguments resolves the same way: configured CloudClient.organization_id
first, then the parent organization of CloudClient.default_workspace_id, then the
parent organization of the authenticated user's default workspace, then the memberships
below.
The deprecated all_organizations=True alias maps to privilege_scope=ANY.
Why the path matters
The two listing paths differ in completeness, not just speed:
- Organization-scoped (an organization or membership scope was resolved) uses the Config API, which filters by name server-side and paginates, so results are complete and each workspace carries its organization attribution.
- Cross-organization (
privilege_scope=INSTANCE_ADMIN, orANYfor an instance admin) uses the public API, which has neither an organization filter nor a name filter. Name matching happens client-side over every visible workspace, and the responses carry no organization attribution.
Searching organizations
Organization search and limits are also server-side. CloudClient.list_organizations
uses the Config API whenever name_contains or limit is passed, which filters and
paginates on the server; with neither argument it uses the public API, which returns
every visible organization in a single request. CloudClient.get_organization fetches
one organization by ID directly, and searches by name through the Config API. These
organization lookup paths fall back to the public listing when the Config API is
unavailable, so self-managed deployments keep working.
1# Copyright (c) 2024 Airbyte, Inc., all rights reserved. 2"""PyAirbyte Cloud client. 3 4## Organization and workspace resolution 5 6Most operations need an organization and/or a workspace context. `CloudClient` derives 7that context from what the caller supplies, falling back to the credentials' own scope, 8the authenticated user's default workspace, and finally to the authenticated user's 9organization memberships. 10 11Two rules govern the whole protocol: an explicitly passed ID always beats an ambient 12one, and within each of those tiers an organization ID beats a workspace ID. 13 14### Where a workspace ID comes from 15 16A workspace ID reaches the client from one of four places, in order: 17 181. The `workspace_id` argument on the operation, such as `CloudClient.get_workspace`. 192. The `X-Airbyte-Workspace-Id` header, when running as an MCP server over HTTP. 203. The `AIRBYTE_CLOUD_WORKSPACE_ID` (or `AIRBYTE_WORKSPACE_ID`) environment variable, 21 read when the client is built with `CloudClient.from_auth(env_vars=True)`. 224. The authenticated user's default workspace from their Airbyte user record. 23 24The configured workspace becomes `CloudClient.default_workspace_id`, the ambient 25workspace context for the client. Workspace-scoped operations use it, then the 26authenticated user's default workspace, whenever no `workspace_id` argument is passed. 27`CloudClient.get_workspace` raises when neither is available. 28 29### If a workspace ID is known 30 31The workspace determines the organization: its parent organization is fetched in a 32single call and used as the organization context. That covers the common case, and 33nothing below applies. 34 35The one exception is an explicit `organization_id` or `organization_name` argument, 36which always wins over a workspace-derived organization — as does a configured 37`CloudClient.organization_id` over an *ambient* workspace. 38 39### If an organization ID is known but no workspace ID 40 41The organization is used as-is, whether it came from the `organization_id` argument, 42from `organization_name` (an exact-name lookup, so it is never used to infer a 43default), or from the credentials as `CloudClient.organization_id`. 44 45### If neither is known 46 47`CloudClient.list_workspaces` uses `privilege_scope` to choose the search breadth: 48`MEMBER_OF` lists direct workspace grants, `ORGANIZATION_ADMIN` lists workspaces in 49member organizations, `INSTANCE_ADMIN` lists every workspace for instance admins, and 50`ANY` chooses the broadest scope available to the caller. `CloudClient.get_organization` 51called with no arguments resolves the same way: configured `CloudClient.organization_id` 52first, then the parent organization of `CloudClient.default_workspace_id`, then the 53parent organization of the authenticated user's default workspace, then the memberships 54below. 55 56The deprecated `all_organizations=True` alias maps to `privilege_scope=ANY`. 57 58### Why the path matters 59 60The two listing paths differ in completeness, not just speed: 61 62- **Organization-scoped** (an organization or membership scope was resolved) uses the 63 Config API, which filters by name server-side and paginates, so results are complete 64 and each workspace carries its organization attribution. 65- **Cross-organization** (`privilege_scope=INSTANCE_ADMIN`, or `ANY` for an instance admin) uses 66 the public API, which has neither an organization filter nor a name filter. Name 67 matching happens client-side over every visible workspace, and the responses carry no 68 organization attribution. 69 70### Searching organizations 71 72Organization search and limits are also server-side. `CloudClient.list_organizations` 73uses the Config API whenever `name_contains` or `limit` is passed, which filters and 74paginates on the server; with neither argument it uses the public API, which returns 75every visible organization in a single request. `CloudClient.get_organization` fetches 76one organization by ID directly, and searches by name through the Config API. These 77organization lookup paths fall back to the public listing when the Config API is 78unavailable, so self-managed deployments keep working. 79""" 80 81from __future__ import annotations 82 83import warnings 84from dataclasses import dataclass, field 85from http import HTTPStatus 86from typing import TYPE_CHECKING, Any, Literal, NoReturn, overload 87 88from airbyte import exceptions as exc 89from airbyte._util import api_util 90from airbyte.cloud._credentials import _AirbyteCredentials 91from airbyte.cloud.models import ( 92 CloudDefaultContextInfo, 93 CloudDefaultWorkspaceUpdateInfo, 94 CloudOrganizationInfo, 95 CloudWorkspaceInfo, 96 WorkspacePrivilegeScope, 97) 98from airbyte.cloud.organizations import CloudOrganization 99from airbyte.cloud.workspaces import CloudWorkspace 100from airbyte.exceptions import AirbyteError, AirbyteMissingResourceError 101 102 103if TYPE_CHECKING: 104 from collections.abc import Callable, Mapping 105 106 from airbyte.secrets.base import SecretString 107 108 109MAX_ORGANIZATION_CANDIDATES = 10 110MAX_WORKSPACES_TO_VALIDATE = 25 111 112 113@dataclass(init=False, kw_only=True) 114class CloudClient: 115 """Authenticated client for Airbyte Cloud and self-managed Airbyte APIs.""" 116 117 _credentials: _AirbyteCredentials 118 _membership_organization_ids: tuple[str, ...] | None 119 _user_permissions: tuple[dict[str, Any], ...] | None 120 _direct_workspace_infos: dict[str, CloudWorkspaceInfo | None] 121 _workspace_organizations: dict[str, CloudOrganizationInfo | None] 122 _validated_direct_workspace_result: tuple[list[CloudWorkspaceInfo], int] | None 123 _authenticated_user_info: dict[str, Any] | None = field(repr=False) 124 _authenticated_user_id: str | None = field(repr=False) 125 _authenticated_bearer_token: SecretString | None 126 127 def __init__( 128 self, 129 *, 130 client_id: str | SecretString | None = None, 131 client_secret: str | SecretString | None = None, 132 bearer_token: str | SecretString | None = None, 133 public_api_root: str | None = None, 134 config_api_root: str | None = None, 135 workspace_id: str | None = None, 136 organization_id: str | None = None, 137 ) -> None: 138 """Initialize a `CloudClient` from explicit auth values.""" 139 self._credentials = _AirbyteCredentials.from_auth( 140 client_id=client_id, 141 client_secret=client_secret, 142 bearer_token=bearer_token, 143 public_api_root=public_api_root, 144 config_api_root=config_api_root, 145 workspace_id=workspace_id, 146 organization_id=organization_id, 147 env_vars=False, 148 ) 149 self._membership_organization_ids = None 150 self._user_permissions = None 151 self._direct_workspace_infos = {} 152 self._workspace_organizations = {} 153 self._validated_direct_workspace_result = None 154 self._authenticated_user_info = None 155 self._authenticated_user_id = None 156 self._authenticated_bearer_token = None 157 158 @property 159 def client_id(self) -> SecretString | None: 160 """OAuth client ID used for authentication.""" 161 return self._credentials.client_id 162 163 @property 164 def client_secret(self) -> SecretString | None: 165 """OAuth client secret used for authentication.""" 166 return self._credentials.client_secret 167 168 @property 169 def bearer_token(self) -> SecretString | None: 170 """Bearer token used for authentication.""" 171 return self._credentials.bearer_token 172 173 @property 174 def public_api_root(self) -> str: 175 """Airbyte Public API root.""" 176 return self._credentials.public_api_root 177 178 @property 179 def config_api_root(self) -> str | None: 180 """Airbyte Config API root.""" 181 return self._credentials.config_api_root 182 183 @property 184 def organization_id(self) -> str | None: 185 """Default organization ID for organization-scoped operations.""" 186 return self._credentials.organization_id 187 188 @property 189 def default_workspace_id(self) -> str | None: 190 """Default workspace ID for workspace-scoped operations.""" 191 return self._credentials.workspace_id 192 193 @classmethod 194 def from_auth( 195 cls, 196 *, 197 env_vars: bool = False, 198 organization_id: str | None = None, 199 client_id: str | SecretString | None = None, 200 client_secret: str | SecretString | None = None, 201 bearer_token: str | SecretString | None = None, 202 public_api_root: str | None = None, 203 config_api_root: str | None = None, 204 ) -> CloudClient: 205 """Create a client from explicit inputs and optionally environment variables. 206 207 When `env_vars` is True, environment variables are checked as a fallback 208 after any explicitly provided values. 209 """ 210 credentials = _AirbyteCredentials.from_auth( 211 organization_id=organization_id, 212 client_id=client_id, 213 client_secret=client_secret, 214 bearer_token=bearer_token, 215 public_api_root=public_api_root, 216 config_api_root=config_api_root, 217 env_vars=env_vars, 218 ) 219 return cls._from_credentials(credentials) 220 221 @classmethod 222 def _from_credentials(cls, credentials: _AirbyteCredentials) -> CloudClient: 223 """Create a client from resolved Cloud credentials.""" 224 return cls( 225 client_id=credentials.client_id, 226 client_secret=credentials.client_secret, 227 bearer_token=credentials.bearer_token, 228 public_api_root=credentials.public_api_root, 229 config_api_root=credentials.config_api_root, 230 workspace_id=credentials.workspace_id, 231 organization_id=credentials.organization_id, 232 ) 233 234 def get_workspace(self, workspace_id: str | None = None) -> CloudWorkspace: 235 """Create a `CloudWorkspace` using this client's credentials. 236 237 See the module docstring for how the workspace is resolved. 238 """ 239 resolved_workspace_id = workspace_id or self.resolve_default_workspace_id() 240 if not resolved_workspace_id: 241 raise exc.PyAirbyteInputError( 242 message="Workspace ID is required.", 243 guidance=( 244 "No workspace was configured, and no default workspace could be resolved " 245 "for the authenticated user. Provide a workspace ID, or call " 246 "`get_default_cloud_context` to discover your workspaces and organizations." 247 ), 248 ) 249 250 credentials = self._credentials.with_workspace_id(resolved_workspace_id) 251 return CloudWorkspace( 252 workspace_id=credentials.workspace_id, 253 client_id=credentials.client_id, 254 client_secret=credentials.client_secret, 255 bearer_token=credentials.bearer_token, 256 api_root=credentials.public_api_root, 257 config_api_root=credentials.config_api_root, 258 ) 259 260 def create_workspace( 261 self, 262 *, 263 name: str, 264 organization_id: str | None = None, 265 region_id: str | None = None, 266 ) -> CloudWorkspaceInfo: 267 """Create an Airbyte workspace.""" 268 resolved_organization_id = organization_id or self.organization_id 269 workspace = api_util.create_workspace( 270 name=name, 271 organization_id=resolved_organization_id, 272 region_id=region_id, 273 api_root=self.public_api_root, 274 client_id=self.client_id, 275 client_secret=self.client_secret, 276 bearer_token=self.bearer_token, 277 ) 278 return CloudWorkspaceInfo.from_api_response(workspace) 279 280 def rename_workspace( 281 self, 282 workspace_id: str, 283 *, 284 name: str, 285 ) -> CloudWorkspaceInfo: 286 """Rename an Airbyte workspace.""" 287 workspace = api_util.rename_workspace( 288 workspace_id=workspace_id, 289 name=name, 290 api_root=self.public_api_root, 291 client_id=self.client_id, 292 client_secret=self.client_secret, 293 bearer_token=self.bearer_token, 294 ) 295 return CloudWorkspaceInfo.from_api_response(workspace) 296 297 def permanently_delete_workspace( 298 self, 299 workspace_id: str, 300 *, 301 workspace_name: str | None = None, 302 safe_mode: bool = True, 303 ) -> None: 304 """Permanently delete an Airbyte workspace if it has no connections. 305 306 When `safe_mode` is enabled, the workspace name must contain `delete-me` 307 or `deleteme`. This also checks for existing connections before deleting 308 and raises `AirbyteWorkspaceNotEmptyError` if the workspace is not empty. 309 """ 310 api_util.permanently_delete_workspace( 311 workspace_id=workspace_id, 312 workspace_name=workspace_name, 313 api_root=self.public_api_root, 314 client_id=self.client_id, 315 client_secret=self.client_secret, 316 bearer_token=self.bearer_token, 317 safe_mode=safe_mode, 318 ) 319 320 @overload 321 def list_workspaces( 322 self, 323 name: str | None = None, 324 *, 325 organization_id: None = None, 326 organization_name: str | None = None, 327 workspace_id: str | None = None, 328 name_contains: str | None = None, 329 name_filter: Callable[[str], bool] | None = None, 330 limit: int | None = None, 331 privilege_scope: WorkspacePrivilegeScope = WorkspacePrivilegeScope.MEMBER_OF, 332 all_organizations: bool = False, 333 ) -> list[CloudWorkspaceInfo]: 334 raise NotImplementedError 335 336 @overload 337 def list_workspaces( 338 self, 339 name: str | None = None, 340 *, 341 organization_id: str, 342 organization_name: str | None = None, 343 workspace_id: str | None = None, 344 name_contains: str | None = None, 345 name_filter: Callable[[str], bool] | None = None, 346 limit: int | None = None, 347 privilege_scope: WorkspacePrivilegeScope = WorkspacePrivilegeScope.MEMBER_OF, 348 all_organizations: bool = False, 349 ) -> list[CloudWorkspaceInfo]: 350 raise NotImplementedError 351 352 def list_workspaces( # noqa: PLR0911, PLR0913 353 self, 354 name: str | None = None, 355 *, 356 organization_id: str | None = None, 357 organization_name: str | None = None, 358 workspace_id: str | None = None, 359 name_contains: str | None = None, 360 name_filter: Callable[[str], bool] | None = None, 361 limit: int | None = None, 362 privilege_scope: WorkspacePrivilegeScope = WorkspacePrivilegeScope.MEMBER_OF, 363 all_organizations: bool = False, 364 ) -> list[CloudWorkspaceInfo]: 365 """List workspaces available to this client. 366 367 `privilege_scope` controls whether this lists direct member workspaces, 368 organization workspaces, or instance-wide workspaces. The deprecated 369 `all_organizations` alias maps to `WorkspacePrivilegeScope.ANY`. 370 """ 371 if limit is not None and limit <= 0: 372 raise exc.PyAirbyteInputError(message="`limit` must be greater than 0.") 373 if organization_id is not None and organization_name is not None: 374 raise exc.PyAirbyteInputError( 375 message="Provide either organization ID or organization name." 376 ) 377 has_explicit_organization = organization_id is not None or organization_name is not None 378 has_explicit_workspace = workspace_id is not None 379 380 if all_organizations: 381 if privilege_scope is not WorkspacePrivilegeScope.MEMBER_OF: 382 raise exc.PyAirbyteInputError( 383 message="all_organizations cannot be combined with privilege_scope." 384 ) 385 warnings.warn( 386 "`all_organizations` is deprecated; use `privilege_scope` instead.", 387 DeprecationWarning, 388 stacklevel=2, 389 ) 390 privilege_scope = WorkspacePrivilegeScope.ANY 391 if name_contains is not None and name_filter is not None: 392 raise exc.PyAirbyteInputError( 393 message="You can provide name_contains or name_filter, but not both." 394 ) 395 if name is not None and name_contains is not None: 396 raise exc.PyAirbyteInputError( 397 message="You can provide name or name_contains, but not both." 398 ) 399 if has_explicit_organization or has_explicit_workspace: 400 resolved_organization_id = self._resolve_workspace_organization_id( 401 organization_id=organization_id, 402 organization_name=organization_name, 403 workspace_id=workspace_id, 404 ) 405 if resolved_organization_id is None: 406 return [] 407 return self._list_workspaces_in_organizations( 408 (resolved_organization_id,), 409 name=name, 410 name_contains=name_contains, 411 name_filter=name_filter, 412 limit=limit, 413 ) 414 415 if privilege_scope is WorkspacePrivilegeScope.MEMBER_OF: 416 return self._list_member_workspaces( 417 name=name, 418 name_contains=name_contains, 419 name_filter=name_filter, 420 limit=limit, 421 ) 422 423 if privilege_scope is WorkspacePrivilegeScope.INSTANCE_ADMIN: 424 if not self._is_instance_admin(): 425 raise exc.PyAirbyteInputError( 426 message="privilege_scope=instance_admin requires the instance_admin permission." 427 ) 428 return self._list_unscoped_workspaces( 429 name=name, 430 name_contains=name_contains, 431 name_filter=name_filter, 432 limit=limit, 433 ) 434 435 if privilege_scope is WorkspacePrivilegeScope.ANY and self._is_instance_admin(): 436 return self._list_unscoped_workspaces( 437 name=name, 438 name_contains=name_contains, 439 name_filter=name_filter, 440 limit=limit, 441 ) 442 443 if privilege_scope in { 444 WorkspacePrivilegeScope.ORGANIZATION_ADMIN, 445 WorkspacePrivilegeScope.ANY, 446 }: 447 organization_ids = self._get_membership_organization_ids() 448 if not organization_ids: 449 return [] 450 return self._list_workspaces_in_organizations( 451 organization_ids, 452 name=name, 453 name_contains=name_contains, 454 name_filter=name_filter, 455 limit=limit, 456 ) 457 458 raise exc.PyAirbyteInputError(message="Unsupported workspace privilege scope.") 459 460 def _list_member_workspaces( 461 self, 462 *, 463 name: str | None = None, 464 name_contains: str | None = None, 465 name_filter: Callable[[str], bool] | None = None, 466 limit: int | None = None, 467 ) -> list[CloudWorkspaceInfo]: 468 """List workspaces granted directly to the authenticated user.""" 469 workspaces, unvalidated_count = self._validate_direct_workspaces() 470 name_substring = name_contains.casefold() if name_contains is not None else None 471 filtered_workspaces: list[CloudWorkspaceInfo] = [] 472 473 def accepts(workspace: CloudWorkspaceInfo) -> bool: 474 if name is not None and workspace.name != name: 475 return False 476 if name_substring is not None and name_substring not in workspace.name.casefold(): 477 return False 478 return name_filter is None or name_filter(workspace.name) 479 480 for workspace in workspaces: 481 if accepts(workspace): 482 filtered_workspaces.append(workspace) 483 if limit is not None and len(filtered_workspaces) == limit: 484 break 485 if unvalidated_count > 0 and (limit is None or len(filtered_workspaces) < limit): 486 for workspace_id in self._get_direct_workspace_ids()[MAX_WORKSPACES_TO_VALIDATE:]: 487 workspace = self._get_direct_workspace_info(workspace_id) 488 if workspace is None or not accepts(workspace): 489 continue 490 filtered_workspaces.append(workspace) 491 if limit is not None and len(filtered_workspaces) == limit: 492 break 493 return filtered_workspaces 494 495 def _list_unscoped_workspaces( 496 self, 497 *, 498 name: str | None, 499 name_contains: str | None, 500 name_filter: Callable[[str], bool] | None, 501 limit: int | None, 502 ) -> list[CloudWorkspaceInfo]: 503 """List workspaces across the instance.""" 504 if name_contains is not None: 505 name_substring = name_contains.casefold() 506 507 def matches_name(workspace_name: str) -> bool: 508 return name_substring in workspace_name.casefold() 509 510 name_filter = matches_name 511 name = None 512 workspaces = api_util.list_workspaces( 513 workspace_id="", 514 api_root=self.public_api_root, 515 client_id=self.client_id, 516 client_secret=self.client_secret, 517 bearer_token=self.bearer_token, 518 name_filter=name_filter, 519 name=name, 520 limit=limit, 521 ) 522 return [CloudWorkspaceInfo.from_api_response(workspace) for workspace in workspaces] 523 524 def _list_workspaces_in_organizations( 525 self, 526 organization_ids: tuple[str, ...], 527 *, 528 name: str | None, 529 name_contains: str | None, 530 name_filter: Callable[[str], bool] | None, 531 limit: int | None, 532 ) -> list[CloudWorkspaceInfo]: 533 """List and combine workspaces from one or more organizations.""" 534 workspace_infos: list[CloudWorkspaceInfo] = [] 535 for organization_id in organization_ids: 536 remaining_limit = None if limit is None else limit - len(workspace_infos) 537 if remaining_limit == 0: 538 break 539 workspaces = api_util.list_workspaces_in_organization( 540 organization_id=organization_id, 541 api_root=self.public_api_root, 542 config_api_root=self.config_api_root, 543 client_id=self.client_id, 544 client_secret=self.client_secret, 545 bearer_token=self._get_config_api_bearer_token(), 546 name_contains=name_contains or name, 547 limit=None if name is not None or name_filter is not None else remaining_limit, 548 ) 549 organization_workspaces = [ 550 CloudWorkspaceInfo.from_mapping(workspace) for workspace in workspaces 551 ] 552 if name is not None: 553 organization_workspaces = [ 554 workspace for workspace in organization_workspaces if workspace.name == name 555 ] 556 if name_filter is not None: 557 organization_workspaces = [ 558 workspace 559 for workspace in organization_workspaces 560 if name_filter(workspace.name) 561 ] 562 workspace_infos.extend(organization_workspaces) 563 if limit is not None and len(workspace_infos) >= limit: 564 break 565 return workspace_infos[:limit] if limit is not None else workspace_infos 566 567 def _resolve_workspace_organization_id( 568 self, 569 *, 570 organization_id: str | None, 571 organization_name: str | None, 572 workspace_id: str | None, 573 ) -> str | None: 574 """Resolve the organization for a workspace listing.""" 575 if organization_id is not None or organization_name is not None: 576 if organization_id is not None: 577 return organization_id 578 # Do not use explicit name lookup to infer a default organization. 579 return self.get_organization(organization_name=organization_name).organization_id 580 581 if workspace_id is not None: 582 return self._get_workspace_parent_organization_id(workspace_id) 583 584 return self._resolve_ambient_organization_id() 585 586 def _resolve_ambient_organization_id(self) -> str | None: 587 """Resolve an organization from configured client context or memberships.""" 588 if self.organization_id is not None: 589 return self.organization_id 590 if self.default_workspace_id is not None: 591 try: 592 return self._get_workspace_parent_organization_id(self.default_workspace_id) 593 except (exc.AirbyteError, exc.PyAirbyteInputError): 594 pass 595 user_default_workspace_id = self._get_user_default_workspace_id() 596 if user_default_workspace_id: 597 try: 598 return self._get_workspace_parent_organization_id(user_default_workspace_id) 599 except (exc.AirbyteError, exc.PyAirbyteInputError): 600 pass 601 602 try: 603 organization_ids = self._get_membership_organization_ids() 604 except (exc.AirbyteError, exc.PyAirbyteInputError): 605 return None 606 if len(organization_ids) > 1: 607 self._raise_ambiguous_organization_error(organization_ids) 608 return organization_ids[0] if organization_ids else None 609 610 def _get_config_api_bearer_token(self) -> SecretString | None: 611 """Get and cache a bearer token for Config API requests.""" 612 if self._authenticated_bearer_token is not None: 613 return self._authenticated_bearer_token 614 if self.bearer_token is not None: 615 self._authenticated_bearer_token = self.bearer_token 616 elif self.client_id is not None and self.client_secret is not None: 617 self._authenticated_bearer_token = api_util.get_bearer_token( 618 client_id=self.client_id, 619 client_secret=self.client_secret, 620 api_root=self.public_api_root, 621 ) 622 return self._authenticated_bearer_token 623 624 def _get_workspace_parent_organization_id(self, workspace_id: str) -> str: 625 """Resolve a workspace's parent organization ID.""" 626 organization = api_util.get_workspace_organization_info( 627 workspace_id=workspace_id, 628 api_root=self.public_api_root, 629 config_api_root=self.config_api_root, 630 client_id=self.client_id, 631 client_secret=self.client_secret, 632 bearer_token=self._get_config_api_bearer_token(), 633 ) 634 resolved_organization_id = organization.get("organizationId") 635 if isinstance(resolved_organization_id, str) and resolved_organization_id: 636 return resolved_organization_id 637 raise exc.PyAirbyteInputError( 638 message="The workspace response did not include an organization ID.", 639 context={"workspace_id": workspace_id, "response": organization}, 640 ) 641 642 def get_workspace_parent_organization_id(self, workspace_id: str) -> str | None: 643 """Return the parent organization ID of a workspace, or `None` if it cannot be resolved.""" 644 try: 645 return self._get_workspace_parent_organization_id(workspace_id) 646 except (exc.AirbyteError, exc.PyAirbyteInputError): 647 return None 648 649 def _get_authenticated_user_info(self) -> dict[str, Any]: 650 """Get and cache the Airbyte user record for the current credentials.""" 651 if self._authenticated_user_info is not None: 652 return self._authenticated_user_info 653 654 bearer_token = self._get_config_api_bearer_token() 655 if bearer_token is None: 656 raise exc.PyAirbyteInputError( 657 message="No authentication credentials provided.", 658 guidance="Provide either client credentials or a bearer token.", 659 ) 660 auth_user_id = api_util.get_user_id_from_bearer_token(bearer_token) 661 self._authenticated_user_info = api_util.get_user_by_auth_id( 662 auth_user_id, 663 api_root=self.public_api_root, 664 config_api_root=self.config_api_root, 665 client_id=self.client_id, 666 client_secret=self.client_secret, 667 bearer_token=bearer_token, 668 ) 669 return self._authenticated_user_info 670 671 def _get_authenticated_user_id(self) -> str: 672 """Get and cache the Airbyte user ID for the current credentials.""" 673 if self._authenticated_user_id is not None: 674 return self._authenticated_user_id 675 676 user = self._get_authenticated_user_info() 677 user_id = user.get("userId") 678 if not isinstance(user_id, str) or not user_id: 679 raise exc.PyAirbyteInputError( 680 message="The Airbyte user response did not include a user ID.", 681 context={"response": user}, 682 ) 683 self._authenticated_user_id = user_id 684 return self._authenticated_user_id 685 686 def _get_user_default_workspace_id(self) -> str | None: 687 """Get the authenticated user's default workspace ID, when available.""" 688 try: 689 default_workspace_id = self._get_authenticated_user_info().get("defaultWorkspaceId") 690 except (exc.AirbyteError, exc.PyAirbyteInputError): 691 return None 692 return ( 693 default_workspace_id 694 if isinstance(default_workspace_id, str) and default_workspace_id 695 else None 696 ) 697 698 def resolve_default_workspace_id(self) -> str | None: 699 """Resolve the configured or authenticated user's default workspace ID.""" 700 configured_workspace_id = self.default_workspace_id 701 if configured_workspace_id: 702 return configured_workspace_id 703 user_default_workspace_id = self._get_user_default_workspace_id() 704 if user_default_workspace_id: 705 return user_default_workspace_id 706 try: 707 live_workspaces, unvalidated_count = self._validate_direct_workspaces() 708 except (AirbyteError, exc.PyAirbyteInputError): 709 return None 710 return ( 711 live_workspaces[0].workspace_id 712 if unvalidated_count == 0 and len(live_workspaces) == 1 713 else None 714 ) 715 716 def _get_user_permissions(self) -> tuple[dict[str, Any], ...]: 717 """Get and cache permissions for the authenticated user.""" 718 if self._user_permissions is None: 719 self._user_permissions = tuple( 720 permission 721 for permission in api_util.list_permissions_for_user( 722 self._get_authenticated_user_id(), 723 api_root=self.public_api_root, 724 config_api_root=self.config_api_root, 725 client_id=self.client_id, 726 client_secret=self.client_secret, 727 bearer_token=self._get_config_api_bearer_token(), 728 ) 729 if isinstance(permission, dict) 730 ) 731 return self._user_permissions 732 733 def _get_membership_organization_ids(self) -> tuple[str, ...]: 734 """Get and cache organization IDs from the caller's permissions.""" 735 if self._membership_organization_ids is not None: 736 return self._membership_organization_ids 737 738 permissions = self._get_user_permissions() 739 organization_ids: list[str] = [] 740 for permission in permissions: 741 permission_organization_id = permission.get("organizationId") 742 if ( 743 isinstance(permission_organization_id, str) 744 and permission_organization_id 745 and permission_organization_id not in organization_ids 746 ): 747 organization_ids.append(permission_organization_id) 748 self._membership_organization_ids = tuple(organization_ids) 749 return self._membership_organization_ids 750 751 def _get_direct_workspace_ids(self) -> tuple[str, ...]: 752 """Get unique workspace IDs from the caller's direct permissions.""" 753 workspace_ids: list[str] = [] 754 for permission in self._get_user_permissions(): 755 workspace_id = permission.get("workspaceId") 756 if isinstance(workspace_id, str) and workspace_id and workspace_id not in workspace_ids: 757 workspace_ids.append(workspace_id) 758 return tuple(workspace_ids) 759 760 def _get_direct_workspace_info(self, workspace_id: str) -> CloudWorkspaceInfo | None: 761 """Fetch a directly granted workspace, or `None` if the grant is stale (404).""" 762 if workspace_id in self._direct_workspace_infos: 763 return self._direct_workspace_infos[workspace_id] 764 try: 765 workspace = api_util.get_workspace( 766 workspace_id=workspace_id, 767 api_root=self.public_api_root, 768 client_id=self.client_id, 769 client_secret=self.client_secret, 770 bearer_token=self.bearer_token, 771 ) 772 except exc.AirbyteMissingResourceError: 773 self._direct_workspace_infos[workspace_id] = None 774 return None 775 workspace_info = CloudWorkspaceInfo.from_api_response(workspace) 776 self._direct_workspace_infos[workspace_id] = workspace_info 777 return workspace_info 778 779 def _get_workspace_organization(self, workspace_id: str) -> CloudOrganizationInfo | None: 780 """Fetch and cache organization info for a workspace.""" 781 if workspace_id in self._workspace_organizations: 782 return self._workspace_organizations[workspace_id] 783 try: 784 organization = api_util.get_workspace_organization_info( 785 workspace_id=workspace_id, 786 api_root=self.public_api_root, 787 config_api_root=self.config_api_root, 788 client_id=self.client_id, 789 client_secret=self.client_secret, 790 bearer_token=self._get_config_api_bearer_token(), 791 ) 792 except (AirbyteError, NotImplementedError): 793 # The workspace is readable via the public API but its organization is not 794 # (e.g. the caller lacks org-level read, or no Config API root can be derived 795 # from a custom public API root). Keep the live workspace and leave the 796 # organization unknown. 797 self._workspace_organizations[workspace_id] = None 798 return None 799 organization_id = organization.get("organizationId") 800 if not isinstance(organization_id, str) or not organization_id: 801 self._workspace_organizations[workspace_id] = None 802 return None 803 organization_info = CloudOrganizationInfo( 804 organization_id=organization_id, 805 organization_name=( 806 organization.get("organizationName") 807 if isinstance(organization.get("organizationName"), str) 808 else None 809 ), 810 ) 811 self._workspace_organizations[workspace_id] = organization_info 812 return organization_info 813 814 def _validate_direct_workspaces(self) -> tuple[list[CloudWorkspaceInfo], int]: 815 """Validate direct workspace grants once within the configured cap.""" 816 if self._validated_direct_workspace_result is not None: 817 return self._validated_direct_workspace_result 818 workspace_ids = self._get_direct_workspace_ids() 819 live_workspaces: list[CloudWorkspaceInfo] = [] 820 for workspace_id in workspace_ids[:MAX_WORKSPACES_TO_VALIDATE]: 821 workspace = self._get_direct_workspace_info(workspace_id) 822 if workspace is None: 823 continue 824 organization = self._get_workspace_organization(workspace_id) 825 if organization is not None: 826 workspace = workspace.model_copy( 827 update={ 828 "organization_id": organization.organization_id, 829 "organization_name": organization.organization_name, 830 } 831 ) 832 live_workspaces.append(workspace) 833 result = ( 834 live_workspaces, 835 max(0, len(workspace_ids) - MAX_WORKSPACES_TO_VALIDATE), 836 ) 837 self._validated_direct_workspace_result = result 838 return result 839 840 def _is_instance_admin(self) -> bool: 841 """Return whether the caller has an instance-admin permission.""" 842 return any( 843 permission.get("permissionType") == "instance_admin" 844 for permission in self._get_user_permissions() 845 ) 846 847 def get_default_context_for_user(self) -> CloudDefaultContextInfo: 848 """Describe the authenticated user's explicit Cloud affinities.""" 849 user_id: str | None = None 850 user_name: str | None = None 851 user_email: str | None = None 852 user: dict[str, Any] | None = None 853 try: 854 user = self._get_authenticated_user_info() 855 except (AirbyteError, exc.PyAirbyteInputError): 856 pass 857 else: 858 user_id = user.get("userId") if isinstance(user.get("userId"), str) else None 859 user_name = user.get("name") if isinstance(user.get("name"), str) else None 860 user_email = user.get("email") if isinstance(user.get("email"), str) else None 861 862 default_workspace_id = self.resolve_default_workspace_id() 863 try: 864 permissions = self._get_user_permissions() 865 except (AirbyteError, exc.PyAirbyteInputError): 866 permissions = () 867 membership_organization_ids = () 868 member_workspaces = [] 869 member_organizations_truncated = False 870 unvalidated_workspace_count = 0 871 else: 872 membership_organization_ids = self._get_membership_organization_ids() 873 member_organizations_truncated = ( 874 len(membership_organization_ids) > MAX_ORGANIZATION_CANDIDATES 875 ) 876 try: 877 member_workspaces, unvalidated_workspace_count = self._validate_direct_workspaces() 878 except (AirbyteError, exc.PyAirbyteInputError): 879 member_workspaces = [] 880 unvalidated_workspace_count = 0 881 member_workspaces = member_workspaces[:] 882 883 member_organizations = [ 884 CloudOrganizationInfo.model_validate(candidate) 885 for candidate in self._get_organization_candidates( 886 membership_organization_ids[:MAX_ORGANIZATION_CANDIDATES] 887 ) 888 ] 889 default_workspace_info: CloudWorkspaceInfo | None = None 890 default_workspace_organization: CloudOrganizationInfo | None = None 891 if default_workspace_id is not None: 892 try: 893 default_workspace_info = self._get_direct_workspace_info(default_workspace_id) 894 except (AirbyteError, exc.PyAirbyteInputError): 895 default_workspace_info = None 896 if default_workspace_info is not None: 897 default_workspace_organization = self._get_workspace_organization( 898 default_workspace_id 899 ) 900 discovery_hints: list[str] = [] 901 if any(permission.get("permissionType") == "instance_admin" for permission in permissions): 902 discovery_hints.append( 903 "Instance-admin access may include every organization and workspace in the " 904 "instance. Use list_cloud_organizations(name_contains=...) or " 905 "list_cloud_workspaces(organization_id=...) to discover others." 906 ) 907 if membership_organization_ids: 908 discovery_hints.append( 909 "Organization membership grants access to every workspace in those " 910 "organizations. Use list_cloud_workspaces(organization_id=<id>) to " 911 "discover workspaces." 912 ) 913 if user is not None and not isinstance(user.get("defaultWorkspaceId"), str): 914 discovery_hints.append( 915 "No default workspace is set. Use " 916 "set_default_cloud_workspace(user_email=<your email>, " 917 "workspace_id=<id>) to durably set one; it applies to both MCP " 918 "sessions and the Airbyte Cloud web app." 919 ) 920 return CloudDefaultContextInfo( 921 user_id=user_id, 922 user_name=user_name, 923 user_email=user_email, 924 default_workspace_id=default_workspace_id, 925 default_workspace_name=( 926 default_workspace_info.name if default_workspace_info else None 927 ), 928 default_workspace_verified=default_workspace_info is not None, 929 default_organization_id=( 930 default_workspace_organization.organization_id 931 if default_workspace_organization is not None 932 else None 933 ), 934 default_organization_name=( 935 default_workspace_organization.organization_name 936 if default_workspace_organization is not None 937 else None 938 ), 939 configured_workspace_id=self.default_workspace_id, 940 configured_organization_id=self.organization_id, 941 member_organizations=member_organizations, 942 member_workspaces=member_workspaces, 943 member_organizations_truncated=member_organizations_truncated, 944 member_workspaces_truncated=unvalidated_workspace_count > 0, 945 unvalidated_workspace_count=unvalidated_workspace_count, 946 discovery_hints=discovery_hints, 947 ) 948 949 def set_default_workspace_for_user( 950 self, 951 *, 952 user_email: str, 953 workspace_id: str, 954 ) -> CloudDefaultWorkspaceUpdateInfo: 955 """Durably set the authenticated user's default workspace. 956 957 This is a persistent, account-level change: it updates the user's stored 958 default workspace in Airbyte Cloud, affecting future sessions and the 959 Airbyte Cloud web app. Validation fails closed: the requested email must 960 match the authenticated user, the workspace must be live (not 961 tombstoned), and the user must be an explicit member of the workspace or 962 its parent organization. 963 """ 964 user = self._get_authenticated_user_info() 965 user_id = user.get("userId") 966 authenticated_email = user.get("email") 967 if not isinstance(user_id, str) or not user_id: 968 raise exc.PyAirbyteInputError( 969 message="The Airbyte user response did not include a user ID.", 970 context={"response": user}, 971 ) 972 if not isinstance(authenticated_email, str) or not authenticated_email: 973 raise exc.PyAirbyteInputError( 974 message="The Airbyte user response did not include an email.", 975 context={"response": user}, 976 ) 977 if user.get("status") == "disabled": 978 raise exc.PyAirbyteInputError( 979 message=( 980 "The authenticated user is disabled/deactivated and cannot update " 981 "a default workspace." 982 ), 983 context={"user_id": user_id, "email": authenticated_email}, 984 ) 985 986 if user_email.strip().lower() != authenticated_email.strip().lower(): 987 raise exc.PyAirbyteInputError( 988 message=("The provided `user_email` does not match the authenticated user."), 989 guidance=( 990 "Call get_default_cloud_context to see the authenticated user " 991 "context, then pass that user's email." 992 ), 993 context={ 994 "user_id": user_id, 995 "email": authenticated_email, 996 "name": user.get("name"), 997 }, 998 ) 999 1000 try: 1001 workspace = api_util.get_workspace_config_api( 1002 workspace_id, 1003 api_root=self.public_api_root, 1004 config_api_root=self.config_api_root, 1005 client_id=self.client_id, 1006 client_secret=self.client_secret, 1007 bearer_token=self._get_config_api_bearer_token(), 1008 ) 1009 except exc.AirbyteMissingResourceError as error: 1010 raise exc.PyAirbyteInputError( 1011 message=f"Workspace {workspace_id} was not found.", 1012 guidance=("Call get_default_cloud_context to see your member workspaces."), 1013 context={"workspace_id": workspace_id}, 1014 ) from error 1015 except exc.AirbyteError as error: 1016 if (error.context or {}).get("status_code") == HTTPStatus.NOT_FOUND: 1017 raise exc.PyAirbyteInputError( 1018 message=f"Workspace {workspace_id} was not found.", 1019 guidance=("Call get_default_cloud_context to see your member workspaces."), 1020 context={"workspace_id": workspace_id}, 1021 ) from error 1022 raise 1023 if workspace.get("tombstone"): 1024 raise exc.PyAirbyteInputError( 1025 message=f"Workspace {workspace_id} is tombstoned (deleted).", 1026 guidance=("Call get_default_cloud_context to see your member workspaces."), 1027 context={"workspace_id": workspace_id}, 1028 ) 1029 workspace_organization_id = workspace.get("organizationId") 1030 if not isinstance(workspace_organization_id, str) or not workspace_organization_id: 1031 workspace_organization_id = None 1032 1033 direct = workspace_id in self._get_direct_workspace_ids() 1034 via_org = ( 1035 workspace_organization_id is not None 1036 and workspace_organization_id in self._get_membership_organization_ids() 1037 ) 1038 if not direct and not via_org: 1039 raise exc.PyAirbyteInputError( 1040 message=( 1041 f"You are not an explicit member of workspace {workspace_id} " 1042 f"or its organization {workspace_organization_id}." 1043 ), 1044 guidance=( 1045 "Call get_default_cloud_context to see your member_workspaces " 1046 "and member_organizations, then choose a workspace you are a " 1047 "member of." 1048 ), 1049 context={ 1050 "workspace_id": workspace_id, 1051 "organization_id": workspace_organization_id, 1052 }, 1053 ) 1054 membership_basis: Literal["workspace", "organization"] = ( 1055 "workspace" if direct else "organization" 1056 ) 1057 1058 if workspace_organization_id is None: 1059 raise exc.PyAirbyteInputError( 1060 message=(f"Workspace {workspace_id} does not belong to an organization."), 1061 guidance=( 1062 "Every Airbyte Cloud workspace must belong to a live " 1063 "organization; the workspace record did not include one." 1064 ), 1065 context={"workspace_id": workspace_id}, 1066 ) 1067 previous_default_workspace_id = user.get("defaultWorkspaceId") 1068 if not isinstance(previous_default_workspace_id, str): 1069 previous_default_workspace_id = None 1070 1071 updated_user = api_util.update_user_default_workspace( 1072 user_id, 1073 workspace_id, 1074 api_root=self.public_api_root, 1075 config_api_root=self.config_api_root, 1076 client_id=self.client_id, 1077 client_secret=self.client_secret, 1078 bearer_token=self._get_config_api_bearer_token(), 1079 ) 1080 if updated_user.get("defaultWorkspaceId") != workspace_id: 1081 raise AirbyteError( 1082 message="Default workspace update did not persist.", 1083 context={ 1084 "workspace_id": workspace_id, 1085 "response": updated_user, 1086 }, 1087 ) 1088 self._authenticated_user_info = None 1089 1090 organization = self._get_workspace_organization(workspace_id) 1091 workspace_name = workspace.get("name") 1092 return CloudDefaultWorkspaceUpdateInfo( 1093 user_id=user_id, 1094 user_email=authenticated_email, 1095 previous_default_workspace_id=previous_default_workspace_id, 1096 default_workspace_id=workspace_id, 1097 default_workspace_name=(workspace_name if isinstance(workspace_name, str) else None), 1098 organization_id=(organization.organization_id if organization is not None else None), 1099 organization_name=( 1100 organization.organization_name if organization is not None else None 1101 ), 1102 membership_basis=membership_basis, 1103 ) 1104 1105 def _get_organization_candidates( 1106 self, 1107 organization_ids: tuple[str, ...], 1108 ) -> list[dict[str, str | None]]: 1109 """Get names for membership-derived organization candidates.""" 1110 candidates: list[dict[str, str | None]] = [] 1111 for organization_id in organization_ids: 1112 organization_name = None 1113 try: 1114 organization_info = api_util.get_organization_info( 1115 organization_id=organization_id, 1116 api_root=self.public_api_root, 1117 config_api_root=self.config_api_root, 1118 client_id=self.client_id, 1119 client_secret=self.client_secret, 1120 bearer_token=self._get_config_api_bearer_token(), 1121 ) 1122 except AirbyteError: 1123 pass 1124 else: 1125 candidate_name = organization_info.get("organizationName") 1126 if isinstance(candidate_name, str): 1127 organization_name = candidate_name 1128 candidates.append( 1129 { 1130 "organization_id": organization_id, 1131 "organization_name": organization_name, 1132 } 1133 ) 1134 return candidates 1135 1136 def _raise_ambiguous_organization_error( 1137 self, 1138 organization_ids: tuple[str, ...], 1139 ) -> NoReturn: 1140 """Raise an error enumerating the caller's candidate organizations.""" 1141 candidates = self._get_organization_candidates( 1142 organization_ids[:MAX_ORGANIZATION_CANDIDATES] 1143 ) 1144 candidate_details = ", ".join( 1145 f"{candidate['organization_id']} " 1146 f"({candidate['organization_name'] or 'name unavailable'})" 1147 for candidate in candidates 1148 ) 1149 raise exc.PyAirbyteInputError( 1150 message=( 1151 "Multiple organization memberships were found for these credentials. Retry " 1152 "with one of these " 1153 "organization IDs " 1154 f"(showing {len(candidates)} of {len(organization_ids)}): {candidate_details}. " 1155 "Call `get_default_cloud_context` to see your memberships." 1156 ), 1157 context={ 1158 "organization_ids": list(organization_ids), 1159 "organization_candidates": candidates, 1160 "total_candidates": len(organization_ids), 1161 }, 1162 ) 1163 1164 def list_organizations( 1165 self, 1166 *, 1167 name_contains: str | None = None, 1168 limit: int | None = None, 1169 ) -> list[CloudOrganization]: 1170 """List organizations available to this client. 1171 1172 See the module docstring for how organization search and limits are resolved. 1173 """ 1174 if limit is not None and limit <= 0: 1175 raise exc.PyAirbyteInputError(message="`limit` must be greater than 0.") 1176 1177 if name_contains is not None or limit is not None: 1178 try: 1179 return self._list_organizations_by_user_id( 1180 name_contains=name_contains, 1181 limit=limit, 1182 ) 1183 except AirbyteError: 1184 pass 1185 1186 organizations = self._fetch_organizations() 1187 if name_contains is not None: 1188 name_substring = name_contains.casefold() 1189 organizations = [ 1190 organization 1191 for organization in organizations 1192 if name_substring in (organization.organization_name or "").casefold() 1193 ] 1194 return organizations if limit is None else organizations[:limit] 1195 1196 def _list_organizations_by_user_id( 1197 self, 1198 *, 1199 name_contains: str | None = None, 1200 limit: int | None = None, 1201 ) -> list[CloudOrganization]: 1202 """List organizations via the Config API, with server-side search and paging.""" 1203 user_id = self._get_authenticated_user_id() 1204 return [ 1205 self._organization_from_mapping(organization) 1206 for organization in api_util.list_organizations_for_user_id( 1207 user_id=user_id, 1208 api_root=self.public_api_root, 1209 config_api_root=self.config_api_root, 1210 client_id=self.client_id, 1211 client_secret=self.client_secret, 1212 bearer_token=self._get_config_api_bearer_token(), 1213 name_contains=name_contains, 1214 limit=limit, 1215 ) 1216 ] 1217 1218 def _organization_from_mapping( 1219 self, 1220 organization: Mapping[str, Any], 1221 ) -> CloudOrganization: 1222 """Build a `CloudOrganization` from a Config API organization mapping.""" 1223 return CloudOrganization( 1224 organization_id=organization["organizationId"], 1225 organization_name=organization.get("organizationName"), 1226 email=organization.get("email"), 1227 client_id=self.client_id, 1228 client_secret=self.client_secret, 1229 bearer_token=self.bearer_token, 1230 public_api_root=self.public_api_root, 1231 config_api_root=self.config_api_root, 1232 ) 1233 1234 def _fetch_organizations(self) -> list[CloudOrganization]: 1235 """Fetch all organizations available to this client.""" 1236 return [ 1237 CloudOrganization( 1238 organization_id=organization.organization_id, 1239 organization_name=organization.organization_name, 1240 email=organization.email, 1241 client_id=self.client_id, 1242 client_secret=self.client_secret, 1243 bearer_token=self.bearer_token, 1244 public_api_root=self.public_api_root, 1245 config_api_root=self.config_api_root, 1246 ) 1247 for organization in api_util.list_organizations_for_user( 1248 api_root=self.public_api_root, 1249 client_id=self.client_id, 1250 client_secret=self.client_secret, 1251 bearer_token=self.bearer_token, 1252 ) 1253 ] 1254 1255 def _resolve_default_organization_id(self) -> str | None: 1256 """Resolve the organization to use when no organization argument is given.""" 1257 return self._resolve_ambient_organization_id() 1258 1259 def _get_organization_by_id(self, organization_id: str) -> CloudOrganization | None: 1260 """Look up a single organization via the Config API, if available.""" 1261 try: 1262 organization_info = api_util.get_organization_info( 1263 organization_id=organization_id, 1264 api_root=self.public_api_root, 1265 config_api_root=self.config_api_root, 1266 client_id=self.client_id, 1267 client_secret=self.client_secret, 1268 bearer_token=self._get_config_api_bearer_token(), 1269 ) 1270 except AirbyteError: 1271 return None 1272 if not isinstance(organization_info.get("organizationId"), str): 1273 return None 1274 return self._organization_from_mapping(organization_info) 1275 1276 def _search_organizations_by_name( 1277 self, 1278 organization_name: str | None, 1279 ) -> list[CloudOrganization]: 1280 """Get organizations whose names contain `organization_name`, if available.""" 1281 if organization_name is not None: 1282 try: 1283 return self._list_organizations_by_user_id(name_contains=organization_name) 1284 except AirbyteError: 1285 pass 1286 return self._fetch_organizations() 1287 1288 def get_organization( 1289 self, 1290 organization_id: str | None = None, 1291 *, 1292 organization_name: str | None = None, 1293 ) -> CloudOrganization: 1294 """Resolve an organization by ID or exact name. 1295 1296 See the module docstring for how the organization is resolved when no 1297 argument is given. 1298 """ 1299 resolved_organization_id = organization_id 1300 if resolved_organization_id and organization_name: 1301 raise exc.PyAirbyteInputError( 1302 message="Provide either organization ID or organization name." 1303 ) 1304 if resolved_organization_id is None and organization_name is None: 1305 resolved_organization_id = self._resolve_default_organization_id() 1306 if not resolved_organization_id and not organization_name: 1307 raise exc.PyAirbyteInputError( 1308 message="Organization ID or organization name is required.", 1309 guidance=( 1310 "Provide an organization ID or name, or call `get_default_cloud_context` " 1311 "to discover your organizations." 1312 ), 1313 ) 1314 1315 if resolved_organization_id: 1316 organization = self._get_organization_by_id(resolved_organization_id) 1317 if organization is not None: 1318 return organization 1319 matching_organizations = [ 1320 candidate 1321 for candidate in self._fetch_organizations() 1322 if candidate.organization_id == resolved_organization_id 1323 ] 1324 else: 1325 matching_organizations = [ 1326 candidate 1327 for candidate in self._search_organizations_by_name(organization_name) 1328 if candidate.organization_name == organization_name 1329 ] 1330 1331 if not matching_organizations: 1332 raise AirbyteMissingResourceError( 1333 resource_type="organization", 1334 resource_name_or_id=resolved_organization_id or organization_name, 1335 ) 1336 if len(matching_organizations) > 1: 1337 total_matches = len(matching_organizations) 1338 shown_matches = matching_organizations[:10] 1339 match_details = ", ".join( 1340 f"{organization.organization_id} ({organization.email or 'email unavailable'})" 1341 for organization in shown_matches 1342 ) 1343 raise exc.PyAirbyteInputError( 1344 message=( 1345 "Organization name matches multiple organizations. Provide an " 1346 f"organization ID to disambiguate. Matching organizations " 1347 f"(showing {len(shown_matches)} of {total_matches}): {match_details}" 1348 ), 1349 context={ 1350 "organization_name": organization_name, 1351 "matching_organizations": [ 1352 { 1353 "organization_id": organization.organization_id, 1354 "email": organization.email, 1355 } 1356 for organization in shown_matches 1357 ], 1358 "total_matches": total_matches, 1359 }, 1360 ) 1361 1362 return matching_organizations[0]
114@dataclass(init=False, kw_only=True) 115class CloudClient: 116 """Authenticated client for Airbyte Cloud and self-managed Airbyte APIs.""" 117 118 _credentials: _AirbyteCredentials 119 _membership_organization_ids: tuple[str, ...] | None 120 _user_permissions: tuple[dict[str, Any], ...] | None 121 _direct_workspace_infos: dict[str, CloudWorkspaceInfo | None] 122 _workspace_organizations: dict[str, CloudOrganizationInfo | None] 123 _validated_direct_workspace_result: tuple[list[CloudWorkspaceInfo], int] | None 124 _authenticated_user_info: dict[str, Any] | None = field(repr=False) 125 _authenticated_user_id: str | None = field(repr=False) 126 _authenticated_bearer_token: SecretString | None 127 128 def __init__( 129 self, 130 *, 131 client_id: str | SecretString | None = None, 132 client_secret: str | SecretString | None = None, 133 bearer_token: str | SecretString | None = None, 134 public_api_root: str | None = None, 135 config_api_root: str | None = None, 136 workspace_id: str | None = None, 137 organization_id: str | None = None, 138 ) -> None: 139 """Initialize a `CloudClient` from explicit auth values.""" 140 self._credentials = _AirbyteCredentials.from_auth( 141 client_id=client_id, 142 client_secret=client_secret, 143 bearer_token=bearer_token, 144 public_api_root=public_api_root, 145 config_api_root=config_api_root, 146 workspace_id=workspace_id, 147 organization_id=organization_id, 148 env_vars=False, 149 ) 150 self._membership_organization_ids = None 151 self._user_permissions = None 152 self._direct_workspace_infos = {} 153 self._workspace_organizations = {} 154 self._validated_direct_workspace_result = None 155 self._authenticated_user_info = None 156 self._authenticated_user_id = None 157 self._authenticated_bearer_token = None 158 159 @property 160 def client_id(self) -> SecretString | None: 161 """OAuth client ID used for authentication.""" 162 return self._credentials.client_id 163 164 @property 165 def client_secret(self) -> SecretString | None: 166 """OAuth client secret used for authentication.""" 167 return self._credentials.client_secret 168 169 @property 170 def bearer_token(self) -> SecretString | None: 171 """Bearer token used for authentication.""" 172 return self._credentials.bearer_token 173 174 @property 175 def public_api_root(self) -> str: 176 """Airbyte Public API root.""" 177 return self._credentials.public_api_root 178 179 @property 180 def config_api_root(self) -> str | None: 181 """Airbyte Config API root.""" 182 return self._credentials.config_api_root 183 184 @property 185 def organization_id(self) -> str | None: 186 """Default organization ID for organization-scoped operations.""" 187 return self._credentials.organization_id 188 189 @property 190 def default_workspace_id(self) -> str | None: 191 """Default workspace ID for workspace-scoped operations.""" 192 return self._credentials.workspace_id 193 194 @classmethod 195 def from_auth( 196 cls, 197 *, 198 env_vars: bool = False, 199 organization_id: str | None = None, 200 client_id: str | SecretString | None = None, 201 client_secret: str | SecretString | None = None, 202 bearer_token: str | SecretString | None = None, 203 public_api_root: str | None = None, 204 config_api_root: str | None = None, 205 ) -> CloudClient: 206 """Create a client from explicit inputs and optionally environment variables. 207 208 When `env_vars` is True, environment variables are checked as a fallback 209 after any explicitly provided values. 210 """ 211 credentials = _AirbyteCredentials.from_auth( 212 organization_id=organization_id, 213 client_id=client_id, 214 client_secret=client_secret, 215 bearer_token=bearer_token, 216 public_api_root=public_api_root, 217 config_api_root=config_api_root, 218 env_vars=env_vars, 219 ) 220 return cls._from_credentials(credentials) 221 222 @classmethod 223 def _from_credentials(cls, credentials: _AirbyteCredentials) -> CloudClient: 224 """Create a client from resolved Cloud credentials.""" 225 return cls( 226 client_id=credentials.client_id, 227 client_secret=credentials.client_secret, 228 bearer_token=credentials.bearer_token, 229 public_api_root=credentials.public_api_root, 230 config_api_root=credentials.config_api_root, 231 workspace_id=credentials.workspace_id, 232 organization_id=credentials.organization_id, 233 ) 234 235 def get_workspace(self, workspace_id: str | None = None) -> CloudWorkspace: 236 """Create a `CloudWorkspace` using this client's credentials. 237 238 See the module docstring for how the workspace is resolved. 239 """ 240 resolved_workspace_id = workspace_id or self.resolve_default_workspace_id() 241 if not resolved_workspace_id: 242 raise exc.PyAirbyteInputError( 243 message="Workspace ID is required.", 244 guidance=( 245 "No workspace was configured, and no default workspace could be resolved " 246 "for the authenticated user. Provide a workspace ID, or call " 247 "`get_default_cloud_context` to discover your workspaces and organizations." 248 ), 249 ) 250 251 credentials = self._credentials.with_workspace_id(resolved_workspace_id) 252 return CloudWorkspace( 253 workspace_id=credentials.workspace_id, 254 client_id=credentials.client_id, 255 client_secret=credentials.client_secret, 256 bearer_token=credentials.bearer_token, 257 api_root=credentials.public_api_root, 258 config_api_root=credentials.config_api_root, 259 ) 260 261 def create_workspace( 262 self, 263 *, 264 name: str, 265 organization_id: str | None = None, 266 region_id: str | None = None, 267 ) -> CloudWorkspaceInfo: 268 """Create an Airbyte workspace.""" 269 resolved_organization_id = organization_id or self.organization_id 270 workspace = api_util.create_workspace( 271 name=name, 272 organization_id=resolved_organization_id, 273 region_id=region_id, 274 api_root=self.public_api_root, 275 client_id=self.client_id, 276 client_secret=self.client_secret, 277 bearer_token=self.bearer_token, 278 ) 279 return CloudWorkspaceInfo.from_api_response(workspace) 280 281 def rename_workspace( 282 self, 283 workspace_id: str, 284 *, 285 name: str, 286 ) -> CloudWorkspaceInfo: 287 """Rename an Airbyte workspace.""" 288 workspace = api_util.rename_workspace( 289 workspace_id=workspace_id, 290 name=name, 291 api_root=self.public_api_root, 292 client_id=self.client_id, 293 client_secret=self.client_secret, 294 bearer_token=self.bearer_token, 295 ) 296 return CloudWorkspaceInfo.from_api_response(workspace) 297 298 def permanently_delete_workspace( 299 self, 300 workspace_id: str, 301 *, 302 workspace_name: str | None = None, 303 safe_mode: bool = True, 304 ) -> None: 305 """Permanently delete an Airbyte workspace if it has no connections. 306 307 When `safe_mode` is enabled, the workspace name must contain `delete-me` 308 or `deleteme`. This also checks for existing connections before deleting 309 and raises `AirbyteWorkspaceNotEmptyError` if the workspace is not empty. 310 """ 311 api_util.permanently_delete_workspace( 312 workspace_id=workspace_id, 313 workspace_name=workspace_name, 314 api_root=self.public_api_root, 315 client_id=self.client_id, 316 client_secret=self.client_secret, 317 bearer_token=self.bearer_token, 318 safe_mode=safe_mode, 319 ) 320 321 @overload 322 def list_workspaces( 323 self, 324 name: str | None = None, 325 *, 326 organization_id: None = None, 327 organization_name: str | None = None, 328 workspace_id: str | None = None, 329 name_contains: str | None = None, 330 name_filter: Callable[[str], bool] | None = None, 331 limit: int | None = None, 332 privilege_scope: WorkspacePrivilegeScope = WorkspacePrivilegeScope.MEMBER_OF, 333 all_organizations: bool = False, 334 ) -> list[CloudWorkspaceInfo]: 335 raise NotImplementedError 336 337 @overload 338 def list_workspaces( 339 self, 340 name: str | None = None, 341 *, 342 organization_id: str, 343 organization_name: str | None = None, 344 workspace_id: str | None = None, 345 name_contains: str | None = None, 346 name_filter: Callable[[str], bool] | None = None, 347 limit: int | None = None, 348 privilege_scope: WorkspacePrivilegeScope = WorkspacePrivilegeScope.MEMBER_OF, 349 all_organizations: bool = False, 350 ) -> list[CloudWorkspaceInfo]: 351 raise NotImplementedError 352 353 def list_workspaces( # noqa: PLR0911, PLR0913 354 self, 355 name: str | None = None, 356 *, 357 organization_id: str | None = None, 358 organization_name: str | None = None, 359 workspace_id: str | None = None, 360 name_contains: str | None = None, 361 name_filter: Callable[[str], bool] | None = None, 362 limit: int | None = None, 363 privilege_scope: WorkspacePrivilegeScope = WorkspacePrivilegeScope.MEMBER_OF, 364 all_organizations: bool = False, 365 ) -> list[CloudWorkspaceInfo]: 366 """List workspaces available to this client. 367 368 `privilege_scope` controls whether this lists direct member workspaces, 369 organization workspaces, or instance-wide workspaces. The deprecated 370 `all_organizations` alias maps to `WorkspacePrivilegeScope.ANY`. 371 """ 372 if limit is not None and limit <= 0: 373 raise exc.PyAirbyteInputError(message="`limit` must be greater than 0.") 374 if organization_id is not None and organization_name is not None: 375 raise exc.PyAirbyteInputError( 376 message="Provide either organization ID or organization name." 377 ) 378 has_explicit_organization = organization_id is not None or organization_name is not None 379 has_explicit_workspace = workspace_id is not None 380 381 if all_organizations: 382 if privilege_scope is not WorkspacePrivilegeScope.MEMBER_OF: 383 raise exc.PyAirbyteInputError( 384 message="all_organizations cannot be combined with privilege_scope." 385 ) 386 warnings.warn( 387 "`all_organizations` is deprecated; use `privilege_scope` instead.", 388 DeprecationWarning, 389 stacklevel=2, 390 ) 391 privilege_scope = WorkspacePrivilegeScope.ANY 392 if name_contains is not None and name_filter is not None: 393 raise exc.PyAirbyteInputError( 394 message="You can provide name_contains or name_filter, but not both." 395 ) 396 if name is not None and name_contains is not None: 397 raise exc.PyAirbyteInputError( 398 message="You can provide name or name_contains, but not both." 399 ) 400 if has_explicit_organization or has_explicit_workspace: 401 resolved_organization_id = self._resolve_workspace_organization_id( 402 organization_id=organization_id, 403 organization_name=organization_name, 404 workspace_id=workspace_id, 405 ) 406 if resolved_organization_id is None: 407 return [] 408 return self._list_workspaces_in_organizations( 409 (resolved_organization_id,), 410 name=name, 411 name_contains=name_contains, 412 name_filter=name_filter, 413 limit=limit, 414 ) 415 416 if privilege_scope is WorkspacePrivilegeScope.MEMBER_OF: 417 return self._list_member_workspaces( 418 name=name, 419 name_contains=name_contains, 420 name_filter=name_filter, 421 limit=limit, 422 ) 423 424 if privilege_scope is WorkspacePrivilegeScope.INSTANCE_ADMIN: 425 if not self._is_instance_admin(): 426 raise exc.PyAirbyteInputError( 427 message="privilege_scope=instance_admin requires the instance_admin permission." 428 ) 429 return self._list_unscoped_workspaces( 430 name=name, 431 name_contains=name_contains, 432 name_filter=name_filter, 433 limit=limit, 434 ) 435 436 if privilege_scope is WorkspacePrivilegeScope.ANY and self._is_instance_admin(): 437 return self._list_unscoped_workspaces( 438 name=name, 439 name_contains=name_contains, 440 name_filter=name_filter, 441 limit=limit, 442 ) 443 444 if privilege_scope in { 445 WorkspacePrivilegeScope.ORGANIZATION_ADMIN, 446 WorkspacePrivilegeScope.ANY, 447 }: 448 organization_ids = self._get_membership_organization_ids() 449 if not organization_ids: 450 return [] 451 return self._list_workspaces_in_organizations( 452 organization_ids, 453 name=name, 454 name_contains=name_contains, 455 name_filter=name_filter, 456 limit=limit, 457 ) 458 459 raise exc.PyAirbyteInputError(message="Unsupported workspace privilege scope.") 460 461 def _list_member_workspaces( 462 self, 463 *, 464 name: str | None = None, 465 name_contains: str | None = None, 466 name_filter: Callable[[str], bool] | None = None, 467 limit: int | None = None, 468 ) -> list[CloudWorkspaceInfo]: 469 """List workspaces granted directly to the authenticated user.""" 470 workspaces, unvalidated_count = self._validate_direct_workspaces() 471 name_substring = name_contains.casefold() if name_contains is not None else None 472 filtered_workspaces: list[CloudWorkspaceInfo] = [] 473 474 def accepts(workspace: CloudWorkspaceInfo) -> bool: 475 if name is not None and workspace.name != name: 476 return False 477 if name_substring is not None and name_substring not in workspace.name.casefold(): 478 return False 479 return name_filter is None or name_filter(workspace.name) 480 481 for workspace in workspaces: 482 if accepts(workspace): 483 filtered_workspaces.append(workspace) 484 if limit is not None and len(filtered_workspaces) == limit: 485 break 486 if unvalidated_count > 0 and (limit is None or len(filtered_workspaces) < limit): 487 for workspace_id in self._get_direct_workspace_ids()[MAX_WORKSPACES_TO_VALIDATE:]: 488 workspace = self._get_direct_workspace_info(workspace_id) 489 if workspace is None or not accepts(workspace): 490 continue 491 filtered_workspaces.append(workspace) 492 if limit is not None and len(filtered_workspaces) == limit: 493 break 494 return filtered_workspaces 495 496 def _list_unscoped_workspaces( 497 self, 498 *, 499 name: str | None, 500 name_contains: str | None, 501 name_filter: Callable[[str], bool] | None, 502 limit: int | None, 503 ) -> list[CloudWorkspaceInfo]: 504 """List workspaces across the instance.""" 505 if name_contains is not None: 506 name_substring = name_contains.casefold() 507 508 def matches_name(workspace_name: str) -> bool: 509 return name_substring in workspace_name.casefold() 510 511 name_filter = matches_name 512 name = None 513 workspaces = api_util.list_workspaces( 514 workspace_id="", 515 api_root=self.public_api_root, 516 client_id=self.client_id, 517 client_secret=self.client_secret, 518 bearer_token=self.bearer_token, 519 name_filter=name_filter, 520 name=name, 521 limit=limit, 522 ) 523 return [CloudWorkspaceInfo.from_api_response(workspace) for workspace in workspaces] 524 525 def _list_workspaces_in_organizations( 526 self, 527 organization_ids: tuple[str, ...], 528 *, 529 name: str | None, 530 name_contains: str | None, 531 name_filter: Callable[[str], bool] | None, 532 limit: int | None, 533 ) -> list[CloudWorkspaceInfo]: 534 """List and combine workspaces from one or more organizations.""" 535 workspace_infos: list[CloudWorkspaceInfo] = [] 536 for organization_id in organization_ids: 537 remaining_limit = None if limit is None else limit - len(workspace_infos) 538 if remaining_limit == 0: 539 break 540 workspaces = api_util.list_workspaces_in_organization( 541 organization_id=organization_id, 542 api_root=self.public_api_root, 543 config_api_root=self.config_api_root, 544 client_id=self.client_id, 545 client_secret=self.client_secret, 546 bearer_token=self._get_config_api_bearer_token(), 547 name_contains=name_contains or name, 548 limit=None if name is not None or name_filter is not None else remaining_limit, 549 ) 550 organization_workspaces = [ 551 CloudWorkspaceInfo.from_mapping(workspace) for workspace in workspaces 552 ] 553 if name is not None: 554 organization_workspaces = [ 555 workspace for workspace in organization_workspaces if workspace.name == name 556 ] 557 if name_filter is not None: 558 organization_workspaces = [ 559 workspace 560 for workspace in organization_workspaces 561 if name_filter(workspace.name) 562 ] 563 workspace_infos.extend(organization_workspaces) 564 if limit is not None and len(workspace_infos) >= limit: 565 break 566 return workspace_infos[:limit] if limit is not None else workspace_infos 567 568 def _resolve_workspace_organization_id( 569 self, 570 *, 571 organization_id: str | None, 572 organization_name: str | None, 573 workspace_id: str | None, 574 ) -> str | None: 575 """Resolve the organization for a workspace listing.""" 576 if organization_id is not None or organization_name is not None: 577 if organization_id is not None: 578 return organization_id 579 # Do not use explicit name lookup to infer a default organization. 580 return self.get_organization(organization_name=organization_name).organization_id 581 582 if workspace_id is not None: 583 return self._get_workspace_parent_organization_id(workspace_id) 584 585 return self._resolve_ambient_organization_id() 586 587 def _resolve_ambient_organization_id(self) -> str | None: 588 """Resolve an organization from configured client context or memberships.""" 589 if self.organization_id is not None: 590 return self.organization_id 591 if self.default_workspace_id is not None: 592 try: 593 return self._get_workspace_parent_organization_id(self.default_workspace_id) 594 except (exc.AirbyteError, exc.PyAirbyteInputError): 595 pass 596 user_default_workspace_id = self._get_user_default_workspace_id() 597 if user_default_workspace_id: 598 try: 599 return self._get_workspace_parent_organization_id(user_default_workspace_id) 600 except (exc.AirbyteError, exc.PyAirbyteInputError): 601 pass 602 603 try: 604 organization_ids = self._get_membership_organization_ids() 605 except (exc.AirbyteError, exc.PyAirbyteInputError): 606 return None 607 if len(organization_ids) > 1: 608 self._raise_ambiguous_organization_error(organization_ids) 609 return organization_ids[0] if organization_ids else None 610 611 def _get_config_api_bearer_token(self) -> SecretString | None: 612 """Get and cache a bearer token for Config API requests.""" 613 if self._authenticated_bearer_token is not None: 614 return self._authenticated_bearer_token 615 if self.bearer_token is not None: 616 self._authenticated_bearer_token = self.bearer_token 617 elif self.client_id is not None and self.client_secret is not None: 618 self._authenticated_bearer_token = api_util.get_bearer_token( 619 client_id=self.client_id, 620 client_secret=self.client_secret, 621 api_root=self.public_api_root, 622 ) 623 return self._authenticated_bearer_token 624 625 def _get_workspace_parent_organization_id(self, workspace_id: str) -> str: 626 """Resolve a workspace's parent organization ID.""" 627 organization = api_util.get_workspace_organization_info( 628 workspace_id=workspace_id, 629 api_root=self.public_api_root, 630 config_api_root=self.config_api_root, 631 client_id=self.client_id, 632 client_secret=self.client_secret, 633 bearer_token=self._get_config_api_bearer_token(), 634 ) 635 resolved_organization_id = organization.get("organizationId") 636 if isinstance(resolved_organization_id, str) and resolved_organization_id: 637 return resolved_organization_id 638 raise exc.PyAirbyteInputError( 639 message="The workspace response did not include an organization ID.", 640 context={"workspace_id": workspace_id, "response": organization}, 641 ) 642 643 def get_workspace_parent_organization_id(self, workspace_id: str) -> str | None: 644 """Return the parent organization ID of a workspace, or `None` if it cannot be resolved.""" 645 try: 646 return self._get_workspace_parent_organization_id(workspace_id) 647 except (exc.AirbyteError, exc.PyAirbyteInputError): 648 return None 649 650 def _get_authenticated_user_info(self) -> dict[str, Any]: 651 """Get and cache the Airbyte user record for the current credentials.""" 652 if self._authenticated_user_info is not None: 653 return self._authenticated_user_info 654 655 bearer_token = self._get_config_api_bearer_token() 656 if bearer_token is None: 657 raise exc.PyAirbyteInputError( 658 message="No authentication credentials provided.", 659 guidance="Provide either client credentials or a bearer token.", 660 ) 661 auth_user_id = api_util.get_user_id_from_bearer_token(bearer_token) 662 self._authenticated_user_info = api_util.get_user_by_auth_id( 663 auth_user_id, 664 api_root=self.public_api_root, 665 config_api_root=self.config_api_root, 666 client_id=self.client_id, 667 client_secret=self.client_secret, 668 bearer_token=bearer_token, 669 ) 670 return self._authenticated_user_info 671 672 def _get_authenticated_user_id(self) -> str: 673 """Get and cache the Airbyte user ID for the current credentials.""" 674 if self._authenticated_user_id is not None: 675 return self._authenticated_user_id 676 677 user = self._get_authenticated_user_info() 678 user_id = user.get("userId") 679 if not isinstance(user_id, str) or not user_id: 680 raise exc.PyAirbyteInputError( 681 message="The Airbyte user response did not include a user ID.", 682 context={"response": user}, 683 ) 684 self._authenticated_user_id = user_id 685 return self._authenticated_user_id 686 687 def _get_user_default_workspace_id(self) -> str | None: 688 """Get the authenticated user's default workspace ID, when available.""" 689 try: 690 default_workspace_id = self._get_authenticated_user_info().get("defaultWorkspaceId") 691 except (exc.AirbyteError, exc.PyAirbyteInputError): 692 return None 693 return ( 694 default_workspace_id 695 if isinstance(default_workspace_id, str) and default_workspace_id 696 else None 697 ) 698 699 def resolve_default_workspace_id(self) -> str | None: 700 """Resolve the configured or authenticated user's default workspace ID.""" 701 configured_workspace_id = self.default_workspace_id 702 if configured_workspace_id: 703 return configured_workspace_id 704 user_default_workspace_id = self._get_user_default_workspace_id() 705 if user_default_workspace_id: 706 return user_default_workspace_id 707 try: 708 live_workspaces, unvalidated_count = self._validate_direct_workspaces() 709 except (AirbyteError, exc.PyAirbyteInputError): 710 return None 711 return ( 712 live_workspaces[0].workspace_id 713 if unvalidated_count == 0 and len(live_workspaces) == 1 714 else None 715 ) 716 717 def _get_user_permissions(self) -> tuple[dict[str, Any], ...]: 718 """Get and cache permissions for the authenticated user.""" 719 if self._user_permissions is None: 720 self._user_permissions = tuple( 721 permission 722 for permission in api_util.list_permissions_for_user( 723 self._get_authenticated_user_id(), 724 api_root=self.public_api_root, 725 config_api_root=self.config_api_root, 726 client_id=self.client_id, 727 client_secret=self.client_secret, 728 bearer_token=self._get_config_api_bearer_token(), 729 ) 730 if isinstance(permission, dict) 731 ) 732 return self._user_permissions 733 734 def _get_membership_organization_ids(self) -> tuple[str, ...]: 735 """Get and cache organization IDs from the caller's permissions.""" 736 if self._membership_organization_ids is not None: 737 return self._membership_organization_ids 738 739 permissions = self._get_user_permissions() 740 organization_ids: list[str] = [] 741 for permission in permissions: 742 permission_organization_id = permission.get("organizationId") 743 if ( 744 isinstance(permission_organization_id, str) 745 and permission_organization_id 746 and permission_organization_id not in organization_ids 747 ): 748 organization_ids.append(permission_organization_id) 749 self._membership_organization_ids = tuple(organization_ids) 750 return self._membership_organization_ids 751 752 def _get_direct_workspace_ids(self) -> tuple[str, ...]: 753 """Get unique workspace IDs from the caller's direct permissions.""" 754 workspace_ids: list[str] = [] 755 for permission in self._get_user_permissions(): 756 workspace_id = permission.get("workspaceId") 757 if isinstance(workspace_id, str) and workspace_id and workspace_id not in workspace_ids: 758 workspace_ids.append(workspace_id) 759 return tuple(workspace_ids) 760 761 def _get_direct_workspace_info(self, workspace_id: str) -> CloudWorkspaceInfo | None: 762 """Fetch a directly granted workspace, or `None` if the grant is stale (404).""" 763 if workspace_id in self._direct_workspace_infos: 764 return self._direct_workspace_infos[workspace_id] 765 try: 766 workspace = api_util.get_workspace( 767 workspace_id=workspace_id, 768 api_root=self.public_api_root, 769 client_id=self.client_id, 770 client_secret=self.client_secret, 771 bearer_token=self.bearer_token, 772 ) 773 except exc.AirbyteMissingResourceError: 774 self._direct_workspace_infos[workspace_id] = None 775 return None 776 workspace_info = CloudWorkspaceInfo.from_api_response(workspace) 777 self._direct_workspace_infos[workspace_id] = workspace_info 778 return workspace_info 779 780 def _get_workspace_organization(self, workspace_id: str) -> CloudOrganizationInfo | None: 781 """Fetch and cache organization info for a workspace.""" 782 if workspace_id in self._workspace_organizations: 783 return self._workspace_organizations[workspace_id] 784 try: 785 organization = api_util.get_workspace_organization_info( 786 workspace_id=workspace_id, 787 api_root=self.public_api_root, 788 config_api_root=self.config_api_root, 789 client_id=self.client_id, 790 client_secret=self.client_secret, 791 bearer_token=self._get_config_api_bearer_token(), 792 ) 793 except (AirbyteError, NotImplementedError): 794 # The workspace is readable via the public API but its organization is not 795 # (e.g. the caller lacks org-level read, or no Config API root can be derived 796 # from a custom public API root). Keep the live workspace and leave the 797 # organization unknown. 798 self._workspace_organizations[workspace_id] = None 799 return None 800 organization_id = organization.get("organizationId") 801 if not isinstance(organization_id, str) or not organization_id: 802 self._workspace_organizations[workspace_id] = None 803 return None 804 organization_info = CloudOrganizationInfo( 805 organization_id=organization_id, 806 organization_name=( 807 organization.get("organizationName") 808 if isinstance(organization.get("organizationName"), str) 809 else None 810 ), 811 ) 812 self._workspace_organizations[workspace_id] = organization_info 813 return organization_info 814 815 def _validate_direct_workspaces(self) -> tuple[list[CloudWorkspaceInfo], int]: 816 """Validate direct workspace grants once within the configured cap.""" 817 if self._validated_direct_workspace_result is not None: 818 return self._validated_direct_workspace_result 819 workspace_ids = self._get_direct_workspace_ids() 820 live_workspaces: list[CloudWorkspaceInfo] = [] 821 for workspace_id in workspace_ids[:MAX_WORKSPACES_TO_VALIDATE]: 822 workspace = self._get_direct_workspace_info(workspace_id) 823 if workspace is None: 824 continue 825 organization = self._get_workspace_organization(workspace_id) 826 if organization is not None: 827 workspace = workspace.model_copy( 828 update={ 829 "organization_id": organization.organization_id, 830 "organization_name": organization.organization_name, 831 } 832 ) 833 live_workspaces.append(workspace) 834 result = ( 835 live_workspaces, 836 max(0, len(workspace_ids) - MAX_WORKSPACES_TO_VALIDATE), 837 ) 838 self._validated_direct_workspace_result = result 839 return result 840 841 def _is_instance_admin(self) -> bool: 842 """Return whether the caller has an instance-admin permission.""" 843 return any( 844 permission.get("permissionType") == "instance_admin" 845 for permission in self._get_user_permissions() 846 ) 847 848 def get_default_context_for_user(self) -> CloudDefaultContextInfo: 849 """Describe the authenticated user's explicit Cloud affinities.""" 850 user_id: str | None = None 851 user_name: str | None = None 852 user_email: str | None = None 853 user: dict[str, Any] | None = None 854 try: 855 user = self._get_authenticated_user_info() 856 except (AirbyteError, exc.PyAirbyteInputError): 857 pass 858 else: 859 user_id = user.get("userId") if isinstance(user.get("userId"), str) else None 860 user_name = user.get("name") if isinstance(user.get("name"), str) else None 861 user_email = user.get("email") if isinstance(user.get("email"), str) else None 862 863 default_workspace_id = self.resolve_default_workspace_id() 864 try: 865 permissions = self._get_user_permissions() 866 except (AirbyteError, exc.PyAirbyteInputError): 867 permissions = () 868 membership_organization_ids = () 869 member_workspaces = [] 870 member_organizations_truncated = False 871 unvalidated_workspace_count = 0 872 else: 873 membership_organization_ids = self._get_membership_organization_ids() 874 member_organizations_truncated = ( 875 len(membership_organization_ids) > MAX_ORGANIZATION_CANDIDATES 876 ) 877 try: 878 member_workspaces, unvalidated_workspace_count = self._validate_direct_workspaces() 879 except (AirbyteError, exc.PyAirbyteInputError): 880 member_workspaces = [] 881 unvalidated_workspace_count = 0 882 member_workspaces = member_workspaces[:] 883 884 member_organizations = [ 885 CloudOrganizationInfo.model_validate(candidate) 886 for candidate in self._get_organization_candidates( 887 membership_organization_ids[:MAX_ORGANIZATION_CANDIDATES] 888 ) 889 ] 890 default_workspace_info: CloudWorkspaceInfo | None = None 891 default_workspace_organization: CloudOrganizationInfo | None = None 892 if default_workspace_id is not None: 893 try: 894 default_workspace_info = self._get_direct_workspace_info(default_workspace_id) 895 except (AirbyteError, exc.PyAirbyteInputError): 896 default_workspace_info = None 897 if default_workspace_info is not None: 898 default_workspace_organization = self._get_workspace_organization( 899 default_workspace_id 900 ) 901 discovery_hints: list[str] = [] 902 if any(permission.get("permissionType") == "instance_admin" for permission in permissions): 903 discovery_hints.append( 904 "Instance-admin access may include every organization and workspace in the " 905 "instance. Use list_cloud_organizations(name_contains=...) or " 906 "list_cloud_workspaces(organization_id=...) to discover others." 907 ) 908 if membership_organization_ids: 909 discovery_hints.append( 910 "Organization membership grants access to every workspace in those " 911 "organizations. Use list_cloud_workspaces(organization_id=<id>) to " 912 "discover workspaces." 913 ) 914 if user is not None and not isinstance(user.get("defaultWorkspaceId"), str): 915 discovery_hints.append( 916 "No default workspace is set. Use " 917 "set_default_cloud_workspace(user_email=<your email>, " 918 "workspace_id=<id>) to durably set one; it applies to both MCP " 919 "sessions and the Airbyte Cloud web app." 920 ) 921 return CloudDefaultContextInfo( 922 user_id=user_id, 923 user_name=user_name, 924 user_email=user_email, 925 default_workspace_id=default_workspace_id, 926 default_workspace_name=( 927 default_workspace_info.name if default_workspace_info else None 928 ), 929 default_workspace_verified=default_workspace_info is not None, 930 default_organization_id=( 931 default_workspace_organization.organization_id 932 if default_workspace_organization is not None 933 else None 934 ), 935 default_organization_name=( 936 default_workspace_organization.organization_name 937 if default_workspace_organization is not None 938 else None 939 ), 940 configured_workspace_id=self.default_workspace_id, 941 configured_organization_id=self.organization_id, 942 member_organizations=member_organizations, 943 member_workspaces=member_workspaces, 944 member_organizations_truncated=member_organizations_truncated, 945 member_workspaces_truncated=unvalidated_workspace_count > 0, 946 unvalidated_workspace_count=unvalidated_workspace_count, 947 discovery_hints=discovery_hints, 948 ) 949 950 def set_default_workspace_for_user( 951 self, 952 *, 953 user_email: str, 954 workspace_id: str, 955 ) -> CloudDefaultWorkspaceUpdateInfo: 956 """Durably set the authenticated user's default workspace. 957 958 This is a persistent, account-level change: it updates the user's stored 959 default workspace in Airbyte Cloud, affecting future sessions and the 960 Airbyte Cloud web app. Validation fails closed: the requested email must 961 match the authenticated user, the workspace must be live (not 962 tombstoned), and the user must be an explicit member of the workspace or 963 its parent organization. 964 """ 965 user = self._get_authenticated_user_info() 966 user_id = user.get("userId") 967 authenticated_email = user.get("email") 968 if not isinstance(user_id, str) or not user_id: 969 raise exc.PyAirbyteInputError( 970 message="The Airbyte user response did not include a user ID.", 971 context={"response": user}, 972 ) 973 if not isinstance(authenticated_email, str) or not authenticated_email: 974 raise exc.PyAirbyteInputError( 975 message="The Airbyte user response did not include an email.", 976 context={"response": user}, 977 ) 978 if user.get("status") == "disabled": 979 raise exc.PyAirbyteInputError( 980 message=( 981 "The authenticated user is disabled/deactivated and cannot update " 982 "a default workspace." 983 ), 984 context={"user_id": user_id, "email": authenticated_email}, 985 ) 986 987 if user_email.strip().lower() != authenticated_email.strip().lower(): 988 raise exc.PyAirbyteInputError( 989 message=("The provided `user_email` does not match the authenticated user."), 990 guidance=( 991 "Call get_default_cloud_context to see the authenticated user " 992 "context, then pass that user's email." 993 ), 994 context={ 995 "user_id": user_id, 996 "email": authenticated_email, 997 "name": user.get("name"), 998 }, 999 ) 1000 1001 try: 1002 workspace = api_util.get_workspace_config_api( 1003 workspace_id, 1004 api_root=self.public_api_root, 1005 config_api_root=self.config_api_root, 1006 client_id=self.client_id, 1007 client_secret=self.client_secret, 1008 bearer_token=self._get_config_api_bearer_token(), 1009 ) 1010 except exc.AirbyteMissingResourceError as error: 1011 raise exc.PyAirbyteInputError( 1012 message=f"Workspace {workspace_id} was not found.", 1013 guidance=("Call get_default_cloud_context to see your member workspaces."), 1014 context={"workspace_id": workspace_id}, 1015 ) from error 1016 except exc.AirbyteError as error: 1017 if (error.context or {}).get("status_code") == HTTPStatus.NOT_FOUND: 1018 raise exc.PyAirbyteInputError( 1019 message=f"Workspace {workspace_id} was not found.", 1020 guidance=("Call get_default_cloud_context to see your member workspaces."), 1021 context={"workspace_id": workspace_id}, 1022 ) from error 1023 raise 1024 if workspace.get("tombstone"): 1025 raise exc.PyAirbyteInputError( 1026 message=f"Workspace {workspace_id} is tombstoned (deleted).", 1027 guidance=("Call get_default_cloud_context to see your member workspaces."), 1028 context={"workspace_id": workspace_id}, 1029 ) 1030 workspace_organization_id = workspace.get("organizationId") 1031 if not isinstance(workspace_organization_id, str) or not workspace_organization_id: 1032 workspace_organization_id = None 1033 1034 direct = workspace_id in self._get_direct_workspace_ids() 1035 via_org = ( 1036 workspace_organization_id is not None 1037 and workspace_organization_id in self._get_membership_organization_ids() 1038 ) 1039 if not direct and not via_org: 1040 raise exc.PyAirbyteInputError( 1041 message=( 1042 f"You are not an explicit member of workspace {workspace_id} " 1043 f"or its organization {workspace_organization_id}." 1044 ), 1045 guidance=( 1046 "Call get_default_cloud_context to see your member_workspaces " 1047 "and member_organizations, then choose a workspace you are a " 1048 "member of." 1049 ), 1050 context={ 1051 "workspace_id": workspace_id, 1052 "organization_id": workspace_organization_id, 1053 }, 1054 ) 1055 membership_basis: Literal["workspace", "organization"] = ( 1056 "workspace" if direct else "organization" 1057 ) 1058 1059 if workspace_organization_id is None: 1060 raise exc.PyAirbyteInputError( 1061 message=(f"Workspace {workspace_id} does not belong to an organization."), 1062 guidance=( 1063 "Every Airbyte Cloud workspace must belong to a live " 1064 "organization; the workspace record did not include one." 1065 ), 1066 context={"workspace_id": workspace_id}, 1067 ) 1068 previous_default_workspace_id = user.get("defaultWorkspaceId") 1069 if not isinstance(previous_default_workspace_id, str): 1070 previous_default_workspace_id = None 1071 1072 updated_user = api_util.update_user_default_workspace( 1073 user_id, 1074 workspace_id, 1075 api_root=self.public_api_root, 1076 config_api_root=self.config_api_root, 1077 client_id=self.client_id, 1078 client_secret=self.client_secret, 1079 bearer_token=self._get_config_api_bearer_token(), 1080 ) 1081 if updated_user.get("defaultWorkspaceId") != workspace_id: 1082 raise AirbyteError( 1083 message="Default workspace update did not persist.", 1084 context={ 1085 "workspace_id": workspace_id, 1086 "response": updated_user, 1087 }, 1088 ) 1089 self._authenticated_user_info = None 1090 1091 organization = self._get_workspace_organization(workspace_id) 1092 workspace_name = workspace.get("name") 1093 return CloudDefaultWorkspaceUpdateInfo( 1094 user_id=user_id, 1095 user_email=authenticated_email, 1096 previous_default_workspace_id=previous_default_workspace_id, 1097 default_workspace_id=workspace_id, 1098 default_workspace_name=(workspace_name if isinstance(workspace_name, str) else None), 1099 organization_id=(organization.organization_id if organization is not None else None), 1100 organization_name=( 1101 organization.organization_name if organization is not None else None 1102 ), 1103 membership_basis=membership_basis, 1104 ) 1105 1106 def _get_organization_candidates( 1107 self, 1108 organization_ids: tuple[str, ...], 1109 ) -> list[dict[str, str | None]]: 1110 """Get names for membership-derived organization candidates.""" 1111 candidates: list[dict[str, str | None]] = [] 1112 for organization_id in organization_ids: 1113 organization_name = None 1114 try: 1115 organization_info = api_util.get_organization_info( 1116 organization_id=organization_id, 1117 api_root=self.public_api_root, 1118 config_api_root=self.config_api_root, 1119 client_id=self.client_id, 1120 client_secret=self.client_secret, 1121 bearer_token=self._get_config_api_bearer_token(), 1122 ) 1123 except AirbyteError: 1124 pass 1125 else: 1126 candidate_name = organization_info.get("organizationName") 1127 if isinstance(candidate_name, str): 1128 organization_name = candidate_name 1129 candidates.append( 1130 { 1131 "organization_id": organization_id, 1132 "organization_name": organization_name, 1133 } 1134 ) 1135 return candidates 1136 1137 def _raise_ambiguous_organization_error( 1138 self, 1139 organization_ids: tuple[str, ...], 1140 ) -> NoReturn: 1141 """Raise an error enumerating the caller's candidate organizations.""" 1142 candidates = self._get_organization_candidates( 1143 organization_ids[:MAX_ORGANIZATION_CANDIDATES] 1144 ) 1145 candidate_details = ", ".join( 1146 f"{candidate['organization_id']} " 1147 f"({candidate['organization_name'] or 'name unavailable'})" 1148 for candidate in candidates 1149 ) 1150 raise exc.PyAirbyteInputError( 1151 message=( 1152 "Multiple organization memberships were found for these credentials. Retry " 1153 "with one of these " 1154 "organization IDs " 1155 f"(showing {len(candidates)} of {len(organization_ids)}): {candidate_details}. " 1156 "Call `get_default_cloud_context` to see your memberships." 1157 ), 1158 context={ 1159 "organization_ids": list(organization_ids), 1160 "organization_candidates": candidates, 1161 "total_candidates": len(organization_ids), 1162 }, 1163 ) 1164 1165 def list_organizations( 1166 self, 1167 *, 1168 name_contains: str | None = None, 1169 limit: int | None = None, 1170 ) -> list[CloudOrganization]: 1171 """List organizations available to this client. 1172 1173 See the module docstring for how organization search and limits are resolved. 1174 """ 1175 if limit is not None and limit <= 0: 1176 raise exc.PyAirbyteInputError(message="`limit` must be greater than 0.") 1177 1178 if name_contains is not None or limit is not None: 1179 try: 1180 return self._list_organizations_by_user_id( 1181 name_contains=name_contains, 1182 limit=limit, 1183 ) 1184 except AirbyteError: 1185 pass 1186 1187 organizations = self._fetch_organizations() 1188 if name_contains is not None: 1189 name_substring = name_contains.casefold() 1190 organizations = [ 1191 organization 1192 for organization in organizations 1193 if name_substring in (organization.organization_name or "").casefold() 1194 ] 1195 return organizations if limit is None else organizations[:limit] 1196 1197 def _list_organizations_by_user_id( 1198 self, 1199 *, 1200 name_contains: str | None = None, 1201 limit: int | None = None, 1202 ) -> list[CloudOrganization]: 1203 """List organizations via the Config API, with server-side search and paging.""" 1204 user_id = self._get_authenticated_user_id() 1205 return [ 1206 self._organization_from_mapping(organization) 1207 for organization in api_util.list_organizations_for_user_id( 1208 user_id=user_id, 1209 api_root=self.public_api_root, 1210 config_api_root=self.config_api_root, 1211 client_id=self.client_id, 1212 client_secret=self.client_secret, 1213 bearer_token=self._get_config_api_bearer_token(), 1214 name_contains=name_contains, 1215 limit=limit, 1216 ) 1217 ] 1218 1219 def _organization_from_mapping( 1220 self, 1221 organization: Mapping[str, Any], 1222 ) -> CloudOrganization: 1223 """Build a `CloudOrganization` from a Config API organization mapping.""" 1224 return CloudOrganization( 1225 organization_id=organization["organizationId"], 1226 organization_name=organization.get("organizationName"), 1227 email=organization.get("email"), 1228 client_id=self.client_id, 1229 client_secret=self.client_secret, 1230 bearer_token=self.bearer_token, 1231 public_api_root=self.public_api_root, 1232 config_api_root=self.config_api_root, 1233 ) 1234 1235 def _fetch_organizations(self) -> list[CloudOrganization]: 1236 """Fetch all organizations available to this client.""" 1237 return [ 1238 CloudOrganization( 1239 organization_id=organization.organization_id, 1240 organization_name=organization.organization_name, 1241 email=organization.email, 1242 client_id=self.client_id, 1243 client_secret=self.client_secret, 1244 bearer_token=self.bearer_token, 1245 public_api_root=self.public_api_root, 1246 config_api_root=self.config_api_root, 1247 ) 1248 for organization in api_util.list_organizations_for_user( 1249 api_root=self.public_api_root, 1250 client_id=self.client_id, 1251 client_secret=self.client_secret, 1252 bearer_token=self.bearer_token, 1253 ) 1254 ] 1255 1256 def _resolve_default_organization_id(self) -> str | None: 1257 """Resolve the organization to use when no organization argument is given.""" 1258 return self._resolve_ambient_organization_id() 1259 1260 def _get_organization_by_id(self, organization_id: str) -> CloudOrganization | None: 1261 """Look up a single organization via the Config API, if available.""" 1262 try: 1263 organization_info = api_util.get_organization_info( 1264 organization_id=organization_id, 1265 api_root=self.public_api_root, 1266 config_api_root=self.config_api_root, 1267 client_id=self.client_id, 1268 client_secret=self.client_secret, 1269 bearer_token=self._get_config_api_bearer_token(), 1270 ) 1271 except AirbyteError: 1272 return None 1273 if not isinstance(organization_info.get("organizationId"), str): 1274 return None 1275 return self._organization_from_mapping(organization_info) 1276 1277 def _search_organizations_by_name( 1278 self, 1279 organization_name: str | None, 1280 ) -> list[CloudOrganization]: 1281 """Get organizations whose names contain `organization_name`, if available.""" 1282 if organization_name is not None: 1283 try: 1284 return self._list_organizations_by_user_id(name_contains=organization_name) 1285 except AirbyteError: 1286 pass 1287 return self._fetch_organizations() 1288 1289 def get_organization( 1290 self, 1291 organization_id: str | None = None, 1292 *, 1293 organization_name: str | None = None, 1294 ) -> CloudOrganization: 1295 """Resolve an organization by ID or exact name. 1296 1297 See the module docstring for how the organization is resolved when no 1298 argument is given. 1299 """ 1300 resolved_organization_id = organization_id 1301 if resolved_organization_id and organization_name: 1302 raise exc.PyAirbyteInputError( 1303 message="Provide either organization ID or organization name." 1304 ) 1305 if resolved_organization_id is None and organization_name is None: 1306 resolved_organization_id = self._resolve_default_organization_id() 1307 if not resolved_organization_id and not organization_name: 1308 raise exc.PyAirbyteInputError( 1309 message="Organization ID or organization name is required.", 1310 guidance=( 1311 "Provide an organization ID or name, or call `get_default_cloud_context` " 1312 "to discover your organizations." 1313 ), 1314 ) 1315 1316 if resolved_organization_id: 1317 organization = self._get_organization_by_id(resolved_organization_id) 1318 if organization is not None: 1319 return organization 1320 matching_organizations = [ 1321 candidate 1322 for candidate in self._fetch_organizations() 1323 if candidate.organization_id == resolved_organization_id 1324 ] 1325 else: 1326 matching_organizations = [ 1327 candidate 1328 for candidate in self._search_organizations_by_name(organization_name) 1329 if candidate.organization_name == organization_name 1330 ] 1331 1332 if not matching_organizations: 1333 raise AirbyteMissingResourceError( 1334 resource_type="organization", 1335 resource_name_or_id=resolved_organization_id or organization_name, 1336 ) 1337 if len(matching_organizations) > 1: 1338 total_matches = len(matching_organizations) 1339 shown_matches = matching_organizations[:10] 1340 match_details = ", ".join( 1341 f"{organization.organization_id} ({organization.email or 'email unavailable'})" 1342 for organization in shown_matches 1343 ) 1344 raise exc.PyAirbyteInputError( 1345 message=( 1346 "Organization name matches multiple organizations. Provide an " 1347 f"organization ID to disambiguate. Matching organizations " 1348 f"(showing {len(shown_matches)} of {total_matches}): {match_details}" 1349 ), 1350 context={ 1351 "organization_name": organization_name, 1352 "matching_organizations": [ 1353 { 1354 "organization_id": organization.organization_id, 1355 "email": organization.email, 1356 } 1357 for organization in shown_matches 1358 ], 1359 "total_matches": total_matches, 1360 }, 1361 ) 1362 1363 return matching_organizations[0]
Authenticated client for Airbyte Cloud and self-managed Airbyte APIs.
128 def __init__( 129 self, 130 *, 131 client_id: str | SecretString | None = None, 132 client_secret: str | SecretString | None = None, 133 bearer_token: str | SecretString | None = None, 134 public_api_root: str | None = None, 135 config_api_root: str | None = None, 136 workspace_id: str | None = None, 137 organization_id: str | None = None, 138 ) -> None: 139 """Initialize a `CloudClient` from explicit auth values.""" 140 self._credentials = _AirbyteCredentials.from_auth( 141 client_id=client_id, 142 client_secret=client_secret, 143 bearer_token=bearer_token, 144 public_api_root=public_api_root, 145 config_api_root=config_api_root, 146 workspace_id=workspace_id, 147 organization_id=organization_id, 148 env_vars=False, 149 ) 150 self._membership_organization_ids = None 151 self._user_permissions = None 152 self._direct_workspace_infos = {} 153 self._workspace_organizations = {} 154 self._validated_direct_workspace_result = None 155 self._authenticated_user_info = None 156 self._authenticated_user_id = None 157 self._authenticated_bearer_token = None
Initialize a CloudClient from explicit auth values.
159 @property 160 def client_id(self) -> SecretString | None: 161 """OAuth client ID used for authentication.""" 162 return self._credentials.client_id
OAuth client ID used for authentication.
164 @property 165 def client_secret(self) -> SecretString | None: 166 """OAuth client secret used for authentication.""" 167 return self._credentials.client_secret
OAuth client secret used for authentication.
169 @property 170 def bearer_token(self) -> SecretString | None: 171 """Bearer token used for authentication.""" 172 return self._credentials.bearer_token
Bearer token used for authentication.
174 @property 175 def public_api_root(self) -> str: 176 """Airbyte Public API root.""" 177 return self._credentials.public_api_root
Airbyte Public API root.
179 @property 180 def config_api_root(self) -> str | None: 181 """Airbyte Config API root.""" 182 return self._credentials.config_api_root
Airbyte Config API root.
184 @property 185 def organization_id(self) -> str | None: 186 """Default organization ID for organization-scoped operations.""" 187 return self._credentials.organization_id
Default organization ID for organization-scoped operations.
189 @property 190 def default_workspace_id(self) -> str | None: 191 """Default workspace ID for workspace-scoped operations.""" 192 return self._credentials.workspace_id
Default workspace ID for workspace-scoped operations.
194 @classmethod 195 def from_auth( 196 cls, 197 *, 198 env_vars: bool = False, 199 organization_id: str | None = None, 200 client_id: str | SecretString | None = None, 201 client_secret: str | SecretString | None = None, 202 bearer_token: str | SecretString | None = None, 203 public_api_root: str | None = None, 204 config_api_root: str | None = None, 205 ) -> CloudClient: 206 """Create a client from explicit inputs and optionally environment variables. 207 208 When `env_vars` is True, environment variables are checked as a fallback 209 after any explicitly provided values. 210 """ 211 credentials = _AirbyteCredentials.from_auth( 212 organization_id=organization_id, 213 client_id=client_id, 214 client_secret=client_secret, 215 bearer_token=bearer_token, 216 public_api_root=public_api_root, 217 config_api_root=config_api_root, 218 env_vars=env_vars, 219 ) 220 return cls._from_credentials(credentials)
Create a client from explicit inputs and optionally environment variables.
When env_vars is True, environment variables are checked as a fallback
after any explicitly provided values.
235 def get_workspace(self, workspace_id: str | None = None) -> CloudWorkspace: 236 """Create a `CloudWorkspace` using this client's credentials. 237 238 See the module docstring for how the workspace is resolved. 239 """ 240 resolved_workspace_id = workspace_id or self.resolve_default_workspace_id() 241 if not resolved_workspace_id: 242 raise exc.PyAirbyteInputError( 243 message="Workspace ID is required.", 244 guidance=( 245 "No workspace was configured, and no default workspace could be resolved " 246 "for the authenticated user. Provide a workspace ID, or call " 247 "`get_default_cloud_context` to discover your workspaces and organizations." 248 ), 249 ) 250 251 credentials = self._credentials.with_workspace_id(resolved_workspace_id) 252 return CloudWorkspace( 253 workspace_id=credentials.workspace_id, 254 client_id=credentials.client_id, 255 client_secret=credentials.client_secret, 256 bearer_token=credentials.bearer_token, 257 api_root=credentials.public_api_root, 258 config_api_root=credentials.config_api_root, 259 )
Create a CloudWorkspace using this client's credentials.
See the module docstring for how the workspace is resolved.
261 def create_workspace( 262 self, 263 *, 264 name: str, 265 organization_id: str | None = None, 266 region_id: str | None = None, 267 ) -> CloudWorkspaceInfo: 268 """Create an Airbyte workspace.""" 269 resolved_organization_id = organization_id or self.organization_id 270 workspace = api_util.create_workspace( 271 name=name, 272 organization_id=resolved_organization_id, 273 region_id=region_id, 274 api_root=self.public_api_root, 275 client_id=self.client_id, 276 client_secret=self.client_secret, 277 bearer_token=self.bearer_token, 278 ) 279 return CloudWorkspaceInfo.from_api_response(workspace)
Create an Airbyte workspace.
281 def rename_workspace( 282 self, 283 workspace_id: str, 284 *, 285 name: str, 286 ) -> CloudWorkspaceInfo: 287 """Rename an Airbyte workspace.""" 288 workspace = api_util.rename_workspace( 289 workspace_id=workspace_id, 290 name=name, 291 api_root=self.public_api_root, 292 client_id=self.client_id, 293 client_secret=self.client_secret, 294 bearer_token=self.bearer_token, 295 ) 296 return CloudWorkspaceInfo.from_api_response(workspace)
Rename an Airbyte workspace.
298 def permanently_delete_workspace( 299 self, 300 workspace_id: str, 301 *, 302 workspace_name: str | None = None, 303 safe_mode: bool = True, 304 ) -> None: 305 """Permanently delete an Airbyte workspace if it has no connections. 306 307 When `safe_mode` is enabled, the workspace name must contain `delete-me` 308 or `deleteme`. This also checks for existing connections before deleting 309 and raises `AirbyteWorkspaceNotEmptyError` if the workspace is not empty. 310 """ 311 api_util.permanently_delete_workspace( 312 workspace_id=workspace_id, 313 workspace_name=workspace_name, 314 api_root=self.public_api_root, 315 client_id=self.client_id, 316 client_secret=self.client_secret, 317 bearer_token=self.bearer_token, 318 safe_mode=safe_mode, 319 )
Permanently delete an Airbyte workspace if it has no connections.
When safe_mode is enabled, the workspace name must contain delete-me
or deleteme. This also checks for existing connections before deleting
and raises AirbyteWorkspaceNotEmptyError if the workspace is not empty.
353 def list_workspaces( # noqa: PLR0911, PLR0913 354 self, 355 name: str | None = None, 356 *, 357 organization_id: str | None = None, 358 organization_name: str | None = None, 359 workspace_id: str | None = None, 360 name_contains: str | None = None, 361 name_filter: Callable[[str], bool] | None = None, 362 limit: int | None = None, 363 privilege_scope: WorkspacePrivilegeScope = WorkspacePrivilegeScope.MEMBER_OF, 364 all_organizations: bool = False, 365 ) -> list[CloudWorkspaceInfo]: 366 """List workspaces available to this client. 367 368 `privilege_scope` controls whether this lists direct member workspaces, 369 organization workspaces, or instance-wide workspaces. The deprecated 370 `all_organizations` alias maps to `WorkspacePrivilegeScope.ANY`. 371 """ 372 if limit is not None and limit <= 0: 373 raise exc.PyAirbyteInputError(message="`limit` must be greater than 0.") 374 if organization_id is not None and organization_name is not None: 375 raise exc.PyAirbyteInputError( 376 message="Provide either organization ID or organization name." 377 ) 378 has_explicit_organization = organization_id is not None or organization_name is not None 379 has_explicit_workspace = workspace_id is not None 380 381 if all_organizations: 382 if privilege_scope is not WorkspacePrivilegeScope.MEMBER_OF: 383 raise exc.PyAirbyteInputError( 384 message="all_organizations cannot be combined with privilege_scope." 385 ) 386 warnings.warn( 387 "`all_organizations` is deprecated; use `privilege_scope` instead.", 388 DeprecationWarning, 389 stacklevel=2, 390 ) 391 privilege_scope = WorkspacePrivilegeScope.ANY 392 if name_contains is not None and name_filter is not None: 393 raise exc.PyAirbyteInputError( 394 message="You can provide name_contains or name_filter, but not both." 395 ) 396 if name is not None and name_contains is not None: 397 raise exc.PyAirbyteInputError( 398 message="You can provide name or name_contains, but not both." 399 ) 400 if has_explicit_organization or has_explicit_workspace: 401 resolved_organization_id = self._resolve_workspace_organization_id( 402 organization_id=organization_id, 403 organization_name=organization_name, 404 workspace_id=workspace_id, 405 ) 406 if resolved_organization_id is None: 407 return [] 408 return self._list_workspaces_in_organizations( 409 (resolved_organization_id,), 410 name=name, 411 name_contains=name_contains, 412 name_filter=name_filter, 413 limit=limit, 414 ) 415 416 if privilege_scope is WorkspacePrivilegeScope.MEMBER_OF: 417 return self._list_member_workspaces( 418 name=name, 419 name_contains=name_contains, 420 name_filter=name_filter, 421 limit=limit, 422 ) 423 424 if privilege_scope is WorkspacePrivilegeScope.INSTANCE_ADMIN: 425 if not self._is_instance_admin(): 426 raise exc.PyAirbyteInputError( 427 message="privilege_scope=instance_admin requires the instance_admin permission." 428 ) 429 return self._list_unscoped_workspaces( 430 name=name, 431 name_contains=name_contains, 432 name_filter=name_filter, 433 limit=limit, 434 ) 435 436 if privilege_scope is WorkspacePrivilegeScope.ANY and self._is_instance_admin(): 437 return self._list_unscoped_workspaces( 438 name=name, 439 name_contains=name_contains, 440 name_filter=name_filter, 441 limit=limit, 442 ) 443 444 if privilege_scope in { 445 WorkspacePrivilegeScope.ORGANIZATION_ADMIN, 446 WorkspacePrivilegeScope.ANY, 447 }: 448 organization_ids = self._get_membership_organization_ids() 449 if not organization_ids: 450 return [] 451 return self._list_workspaces_in_organizations( 452 organization_ids, 453 name=name, 454 name_contains=name_contains, 455 name_filter=name_filter, 456 limit=limit, 457 ) 458 459 raise exc.PyAirbyteInputError(message="Unsupported workspace privilege scope.")
List workspaces available to this client.
privilege_scope controls whether this lists direct member workspaces,
organization workspaces, or instance-wide workspaces. The deprecated
all_organizations alias maps to WorkspacePrivilegeScope.ANY.
643 def get_workspace_parent_organization_id(self, workspace_id: str) -> str | None: 644 """Return the parent organization ID of a workspace, or `None` if it cannot be resolved.""" 645 try: 646 return self._get_workspace_parent_organization_id(workspace_id) 647 except (exc.AirbyteError, exc.PyAirbyteInputError): 648 return None
Return the parent organization ID of a workspace, or None if it cannot be resolved.
699 def resolve_default_workspace_id(self) -> str | None: 700 """Resolve the configured or authenticated user's default workspace ID.""" 701 configured_workspace_id = self.default_workspace_id 702 if configured_workspace_id: 703 return configured_workspace_id 704 user_default_workspace_id = self._get_user_default_workspace_id() 705 if user_default_workspace_id: 706 return user_default_workspace_id 707 try: 708 live_workspaces, unvalidated_count = self._validate_direct_workspaces() 709 except (AirbyteError, exc.PyAirbyteInputError): 710 return None 711 return ( 712 live_workspaces[0].workspace_id 713 if unvalidated_count == 0 and len(live_workspaces) == 1 714 else None 715 )
Resolve the configured or authenticated user's default workspace ID.
848 def get_default_context_for_user(self) -> CloudDefaultContextInfo: 849 """Describe the authenticated user's explicit Cloud affinities.""" 850 user_id: str | None = None 851 user_name: str | None = None 852 user_email: str | None = None 853 user: dict[str, Any] | None = None 854 try: 855 user = self._get_authenticated_user_info() 856 except (AirbyteError, exc.PyAirbyteInputError): 857 pass 858 else: 859 user_id = user.get("userId") if isinstance(user.get("userId"), str) else None 860 user_name = user.get("name") if isinstance(user.get("name"), str) else None 861 user_email = user.get("email") if isinstance(user.get("email"), str) else None 862 863 default_workspace_id = self.resolve_default_workspace_id() 864 try: 865 permissions = self._get_user_permissions() 866 except (AirbyteError, exc.PyAirbyteInputError): 867 permissions = () 868 membership_organization_ids = () 869 member_workspaces = [] 870 member_organizations_truncated = False 871 unvalidated_workspace_count = 0 872 else: 873 membership_organization_ids = self._get_membership_organization_ids() 874 member_organizations_truncated = ( 875 len(membership_organization_ids) > MAX_ORGANIZATION_CANDIDATES 876 ) 877 try: 878 member_workspaces, unvalidated_workspace_count = self._validate_direct_workspaces() 879 except (AirbyteError, exc.PyAirbyteInputError): 880 member_workspaces = [] 881 unvalidated_workspace_count = 0 882 member_workspaces = member_workspaces[:] 883 884 member_organizations = [ 885 CloudOrganizationInfo.model_validate(candidate) 886 for candidate in self._get_organization_candidates( 887 membership_organization_ids[:MAX_ORGANIZATION_CANDIDATES] 888 ) 889 ] 890 default_workspace_info: CloudWorkspaceInfo | None = None 891 default_workspace_organization: CloudOrganizationInfo | None = None 892 if default_workspace_id is not None: 893 try: 894 default_workspace_info = self._get_direct_workspace_info(default_workspace_id) 895 except (AirbyteError, exc.PyAirbyteInputError): 896 default_workspace_info = None 897 if default_workspace_info is not None: 898 default_workspace_organization = self._get_workspace_organization( 899 default_workspace_id 900 ) 901 discovery_hints: list[str] = [] 902 if any(permission.get("permissionType") == "instance_admin" for permission in permissions): 903 discovery_hints.append( 904 "Instance-admin access may include every organization and workspace in the " 905 "instance. Use list_cloud_organizations(name_contains=...) or " 906 "list_cloud_workspaces(organization_id=...) to discover others." 907 ) 908 if membership_organization_ids: 909 discovery_hints.append( 910 "Organization membership grants access to every workspace in those " 911 "organizations. Use list_cloud_workspaces(organization_id=<id>) to " 912 "discover workspaces." 913 ) 914 if user is not None and not isinstance(user.get("defaultWorkspaceId"), str): 915 discovery_hints.append( 916 "No default workspace is set. Use " 917 "set_default_cloud_workspace(user_email=<your email>, " 918 "workspace_id=<id>) to durably set one; it applies to both MCP " 919 "sessions and the Airbyte Cloud web app." 920 ) 921 return CloudDefaultContextInfo( 922 user_id=user_id, 923 user_name=user_name, 924 user_email=user_email, 925 default_workspace_id=default_workspace_id, 926 default_workspace_name=( 927 default_workspace_info.name if default_workspace_info else None 928 ), 929 default_workspace_verified=default_workspace_info is not None, 930 default_organization_id=( 931 default_workspace_organization.organization_id 932 if default_workspace_organization is not None 933 else None 934 ), 935 default_organization_name=( 936 default_workspace_organization.organization_name 937 if default_workspace_organization is not None 938 else None 939 ), 940 configured_workspace_id=self.default_workspace_id, 941 configured_organization_id=self.organization_id, 942 member_organizations=member_organizations, 943 member_workspaces=member_workspaces, 944 member_organizations_truncated=member_organizations_truncated, 945 member_workspaces_truncated=unvalidated_workspace_count > 0, 946 unvalidated_workspace_count=unvalidated_workspace_count, 947 discovery_hints=discovery_hints, 948 )
Describe the authenticated user's explicit Cloud affinities.
950 def set_default_workspace_for_user( 951 self, 952 *, 953 user_email: str, 954 workspace_id: str, 955 ) -> CloudDefaultWorkspaceUpdateInfo: 956 """Durably set the authenticated user's default workspace. 957 958 This is a persistent, account-level change: it updates the user's stored 959 default workspace in Airbyte Cloud, affecting future sessions and the 960 Airbyte Cloud web app. Validation fails closed: the requested email must 961 match the authenticated user, the workspace must be live (not 962 tombstoned), and the user must be an explicit member of the workspace or 963 its parent organization. 964 """ 965 user = self._get_authenticated_user_info() 966 user_id = user.get("userId") 967 authenticated_email = user.get("email") 968 if not isinstance(user_id, str) or not user_id: 969 raise exc.PyAirbyteInputError( 970 message="The Airbyte user response did not include a user ID.", 971 context={"response": user}, 972 ) 973 if not isinstance(authenticated_email, str) or not authenticated_email: 974 raise exc.PyAirbyteInputError( 975 message="The Airbyte user response did not include an email.", 976 context={"response": user}, 977 ) 978 if user.get("status") == "disabled": 979 raise exc.PyAirbyteInputError( 980 message=( 981 "The authenticated user is disabled/deactivated and cannot update " 982 "a default workspace." 983 ), 984 context={"user_id": user_id, "email": authenticated_email}, 985 ) 986 987 if user_email.strip().lower() != authenticated_email.strip().lower(): 988 raise exc.PyAirbyteInputError( 989 message=("The provided `user_email` does not match the authenticated user."), 990 guidance=( 991 "Call get_default_cloud_context to see the authenticated user " 992 "context, then pass that user's email." 993 ), 994 context={ 995 "user_id": user_id, 996 "email": authenticated_email, 997 "name": user.get("name"), 998 }, 999 ) 1000 1001 try: 1002 workspace = api_util.get_workspace_config_api( 1003 workspace_id, 1004 api_root=self.public_api_root, 1005 config_api_root=self.config_api_root, 1006 client_id=self.client_id, 1007 client_secret=self.client_secret, 1008 bearer_token=self._get_config_api_bearer_token(), 1009 ) 1010 except exc.AirbyteMissingResourceError as error: 1011 raise exc.PyAirbyteInputError( 1012 message=f"Workspace {workspace_id} was not found.", 1013 guidance=("Call get_default_cloud_context to see your member workspaces."), 1014 context={"workspace_id": workspace_id}, 1015 ) from error 1016 except exc.AirbyteError as error: 1017 if (error.context or {}).get("status_code") == HTTPStatus.NOT_FOUND: 1018 raise exc.PyAirbyteInputError( 1019 message=f"Workspace {workspace_id} was not found.", 1020 guidance=("Call get_default_cloud_context to see your member workspaces."), 1021 context={"workspace_id": workspace_id}, 1022 ) from error 1023 raise 1024 if workspace.get("tombstone"): 1025 raise exc.PyAirbyteInputError( 1026 message=f"Workspace {workspace_id} is tombstoned (deleted).", 1027 guidance=("Call get_default_cloud_context to see your member workspaces."), 1028 context={"workspace_id": workspace_id}, 1029 ) 1030 workspace_organization_id = workspace.get("organizationId") 1031 if not isinstance(workspace_organization_id, str) or not workspace_organization_id: 1032 workspace_organization_id = None 1033 1034 direct = workspace_id in self._get_direct_workspace_ids() 1035 via_org = ( 1036 workspace_organization_id is not None 1037 and workspace_organization_id in self._get_membership_organization_ids() 1038 ) 1039 if not direct and not via_org: 1040 raise exc.PyAirbyteInputError( 1041 message=( 1042 f"You are not an explicit member of workspace {workspace_id} " 1043 f"or its organization {workspace_organization_id}." 1044 ), 1045 guidance=( 1046 "Call get_default_cloud_context to see your member_workspaces " 1047 "and member_organizations, then choose a workspace you are a " 1048 "member of." 1049 ), 1050 context={ 1051 "workspace_id": workspace_id, 1052 "organization_id": workspace_organization_id, 1053 }, 1054 ) 1055 membership_basis: Literal["workspace", "organization"] = ( 1056 "workspace" if direct else "organization" 1057 ) 1058 1059 if workspace_organization_id is None: 1060 raise exc.PyAirbyteInputError( 1061 message=(f"Workspace {workspace_id} does not belong to an organization."), 1062 guidance=( 1063 "Every Airbyte Cloud workspace must belong to a live " 1064 "organization; the workspace record did not include one." 1065 ), 1066 context={"workspace_id": workspace_id}, 1067 ) 1068 previous_default_workspace_id = user.get("defaultWorkspaceId") 1069 if not isinstance(previous_default_workspace_id, str): 1070 previous_default_workspace_id = None 1071 1072 updated_user = api_util.update_user_default_workspace( 1073 user_id, 1074 workspace_id, 1075 api_root=self.public_api_root, 1076 config_api_root=self.config_api_root, 1077 client_id=self.client_id, 1078 client_secret=self.client_secret, 1079 bearer_token=self._get_config_api_bearer_token(), 1080 ) 1081 if updated_user.get("defaultWorkspaceId") != workspace_id: 1082 raise AirbyteError( 1083 message="Default workspace update did not persist.", 1084 context={ 1085 "workspace_id": workspace_id, 1086 "response": updated_user, 1087 }, 1088 ) 1089 self._authenticated_user_info = None 1090 1091 organization = self._get_workspace_organization(workspace_id) 1092 workspace_name = workspace.get("name") 1093 return CloudDefaultWorkspaceUpdateInfo( 1094 user_id=user_id, 1095 user_email=authenticated_email, 1096 previous_default_workspace_id=previous_default_workspace_id, 1097 default_workspace_id=workspace_id, 1098 default_workspace_name=(workspace_name if isinstance(workspace_name, str) else None), 1099 organization_id=(organization.organization_id if organization is not None else None), 1100 organization_name=( 1101 organization.organization_name if organization is not None else None 1102 ), 1103 membership_basis=membership_basis, 1104 )
Durably set the authenticated user's default workspace.
This is a persistent, account-level change: it updates the user's stored default workspace in Airbyte Cloud, affecting future sessions and the Airbyte Cloud web app. Validation fails closed: the requested email must match the authenticated user, the workspace must be live (not tombstoned), and the user must be an explicit member of the workspace or its parent organization.
1165 def list_organizations( 1166 self, 1167 *, 1168 name_contains: str | None = None, 1169 limit: int | None = None, 1170 ) -> list[CloudOrganization]: 1171 """List organizations available to this client. 1172 1173 See the module docstring for how organization search and limits are resolved. 1174 """ 1175 if limit is not None and limit <= 0: 1176 raise exc.PyAirbyteInputError(message="`limit` must be greater than 0.") 1177 1178 if name_contains is not None or limit is not None: 1179 try: 1180 return self._list_organizations_by_user_id( 1181 name_contains=name_contains, 1182 limit=limit, 1183 ) 1184 except AirbyteError: 1185 pass 1186 1187 organizations = self._fetch_organizations() 1188 if name_contains is not None: 1189 name_substring = name_contains.casefold() 1190 organizations = [ 1191 organization 1192 for organization in organizations 1193 if name_substring in (organization.organization_name or "").casefold() 1194 ] 1195 return organizations if limit is None else organizations[:limit]
List organizations available to this client.
See the module docstring for how organization search and limits are resolved.
1289 def get_organization( 1290 self, 1291 organization_id: str | None = None, 1292 *, 1293 organization_name: str | None = None, 1294 ) -> CloudOrganization: 1295 """Resolve an organization by ID or exact name. 1296 1297 See the module docstring for how the organization is resolved when no 1298 argument is given. 1299 """ 1300 resolved_organization_id = organization_id 1301 if resolved_organization_id and organization_name: 1302 raise exc.PyAirbyteInputError( 1303 message="Provide either organization ID or organization name." 1304 ) 1305 if resolved_organization_id is None and organization_name is None: 1306 resolved_organization_id = self._resolve_default_organization_id() 1307 if not resolved_organization_id and not organization_name: 1308 raise exc.PyAirbyteInputError( 1309 message="Organization ID or organization name is required.", 1310 guidance=( 1311 "Provide an organization ID or name, or call `get_default_cloud_context` " 1312 "to discover your organizations." 1313 ), 1314 ) 1315 1316 if resolved_organization_id: 1317 organization = self._get_organization_by_id(resolved_organization_id) 1318 if organization is not None: 1319 return organization 1320 matching_organizations = [ 1321 candidate 1322 for candidate in self._fetch_organizations() 1323 if candidate.organization_id == resolved_organization_id 1324 ] 1325 else: 1326 matching_organizations = [ 1327 candidate 1328 for candidate in self._search_organizations_by_name(organization_name) 1329 if candidate.organization_name == organization_name 1330 ] 1331 1332 if not matching_organizations: 1333 raise AirbyteMissingResourceError( 1334 resource_type="organization", 1335 resource_name_or_id=resolved_organization_id or organization_name, 1336 ) 1337 if len(matching_organizations) > 1: 1338 total_matches = len(matching_organizations) 1339 shown_matches = matching_organizations[:10] 1340 match_details = ", ".join( 1341 f"{organization.organization_id} ({organization.email or 'email unavailable'})" 1342 for organization in shown_matches 1343 ) 1344 raise exc.PyAirbyteInputError( 1345 message=( 1346 "Organization name matches multiple organizations. Provide an " 1347 f"organization ID to disambiguate. Matching organizations " 1348 f"(showing {len(shown_matches)} of {total_matches}): {match_details}" 1349 ), 1350 context={ 1351 "organization_name": organization_name, 1352 "matching_organizations": [ 1353 { 1354 "organization_id": organization.organization_id, 1355 "email": organization.email, 1356 } 1357 for organization in shown_matches 1358 ], 1359 "total_matches": total_matches, 1360 }, 1361 ) 1362 1363 return matching_organizations[0]
Resolve an organization by ID or exact name.
See the module docstring for how the organization is resolved when no argument is given.