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:

  1. The workspace_id argument on the operation, such as CloudClient.get_workspace.
  2. The X-Airbyte-Workspace-Id header, when running as an MCP server over HTTP.
  3. The AIRBYTE_CLOUD_WORKSPACE_ID (or AIRBYTE_WORKSPACE_ID) environment variable, read when the client is built with CloudClient.from_auth(env_vars=True).
  4. 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, or ANY for 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]
MAX_ORGANIZATION_CANDIDATES = 10
MAX_WORKSPACES_TO_VALIDATE = 25
@dataclass(init=False, kw_only=True)
class CloudClient:
 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.

CloudClient( *, client_id: str | airbyte.secrets.SecretString | None = None, client_secret: str | airbyte.secrets.SecretString | None = None, bearer_token: str | airbyte.secrets.SecretString | None = None, public_api_root: str | None = None, config_api_root: str | None = None, workspace_id: str | None = None, organization_id: str | None = None)
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.

client_id: airbyte.secrets.SecretString | None
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.

client_secret: airbyte.secrets.SecretString | None
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.

bearer_token: airbyte.secrets.SecretString | None
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.

public_api_root: str
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.

config_api_root: str | None
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.

organization_id: str | None
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.

default_workspace_id: str | None
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.

@classmethod
def from_auth( cls, *, env_vars: bool = False, organization_id: str | None = None, client_id: str | airbyte.secrets.SecretString | None = None, client_secret: str | airbyte.secrets.SecretString | None = None, bearer_token: str | airbyte.secrets.SecretString | None = None, public_api_root: str | None = None, config_api_root: str | None = None) -> CloudClient:
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.

def get_workspace( self, workspace_id: str | None = None) -> airbyte.cloud.CloudWorkspace:
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.

def create_workspace( self, *, name: str, organization_id: str | None = None, region_id: str | None = None) -> airbyte.cloud.CloudWorkspaceInfo:
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.

def rename_workspace( self, workspace_id: str, *, name: str) -> airbyte.cloud.CloudWorkspaceInfo:
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.

def permanently_delete_workspace( self, workspace_id: str, *, workspace_name: str | None = None, safe_mode: bool = True) -> None:
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.

def list_workspaces( self, name: str | None = None, *, organization_id: str | None = None, organization_name: str | None = None, workspace_id: str | None = None, name_contains: str | None = None, name_filter: Callable[[str], bool] | None = None, limit: int | None = None, privilege_scope: airbyte.cloud.WorkspacePrivilegeScope = <WorkspacePrivilegeScope.MEMBER_OF: 'member_of'>, all_organizations: bool = False) -> list[airbyte.cloud.CloudWorkspaceInfo]:
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.

def get_workspace_parent_organization_id(self, workspace_id: str) -> str | None:
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.

def resolve_default_workspace_id(self) -> str | None:
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.

def get_default_context_for_user(self) -> airbyte.cloud.CloudDefaultContextInfo:
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.

def set_default_workspace_for_user( self, *, user_email: str, workspace_id: str) -> airbyte.cloud.models.CloudDefaultWorkspaceUpdateInfo:
 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.

def list_organizations( self, *, name_contains: str | None = None, limit: int | None = None) -> list[airbyte.cloud.CloudOrganization]:
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.

def get_organization( self, organization_id: str | None = None, *, organization_name: str | None = None) -> airbyte.cloud.CloudOrganization:
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.