airbyte.cloud

PyAirbyte classes and methods for interacting with the Airbyte Cloud API.

You can use this module to interact with Airbyte Cloud, OSS, and Enterprise.

Self-managed Airbyte instances

For self-managed Airbyte instances, set api_root to the Public API root for your deployment. For the default self-managed route, that usually ends in /api/public/v1. PyAirbyte uses the Public API for workspace and organization discovery.

Some Cloud module methods also call the Config API, including methods such as CloudConnection.dump_raw_catalog(), which reads the configured catalog directly from Airbyte. For documented self-managed deployments where the Public API root ends in /api/public/v1, PyAirbyte infers the Config API root by replacing that suffix with /api/v1.

If your deployment uses custom ingress or a nonstandard reverse proxy, pass config_api_root explicitly or set the AIRBYTE_CLOUD_CONFIG_API_URL environment variable.

from airbyte import cloud

workspace = cloud.CloudWorkspace(
    workspace_id="...",
    client_id="...",
    client_secret="...",
    api_root="https://airbyte.example.com/api/public/v1",
    config_api_root="https://airbyte.example.com/api/v1",
)

connection = workspace.get_connection(connection_id="...")
raw_catalog = connection.dump_raw_catalog()

Examples

Basic Sync Example:

import airbyte as ab
from airbyte import cloud

# Initialize an Airbyte Cloud workspace object
workspace = cloud.CloudWorkspace(
    workspace_id="123",
    api_key=ab.get_secret("AIRBYTE_CLOUD_API_KEY"),
)

# Run a sync job on Airbyte Cloud
connection = workspace.get_connection(connection_id="456")
sync_result = connection.run_sync()
print(sync_result.get_job_status())

Example Read From Cloud Destination:

If your destination is supported, you can read records directly from the SyncResult object. Currently this is supported in Snowflake and BigQuery only.

# Assuming we've already created a `connection` object...

# Get the latest job result and print the stream names
sync_result = connection.get_sync_result()
print(sync_result.stream_names)

# Get a dataset from the sync result
dataset: CachedDataset = sync_result.get_dataset("users")

# Get a SQLAlchemy table to use in SQL queries...
users_table = dataset.to_sql_table()
print(f"Table name: {users_table.name}")

# Or iterate over the dataset directly
for record in dataset:
    print(record)
  1# Copyright (c) 2024 Airbyte, Inc., all rights reserved.
  2"""PyAirbyte classes and methods for interacting with the Airbyte Cloud API.
  3
  4You can use this module to interact with Airbyte Cloud, OSS, and Enterprise.
  5
  6## Self-managed Airbyte instances
  7
  8For self-managed Airbyte instances, set `api_root` to the Public API root for your
  9deployment. For the default self-managed route, that usually ends in `/api/public/v1`.
 10PyAirbyte uses the Public API for workspace and organization discovery.
 11
 12Some Cloud module methods also call the Config API, including methods such as
 13`CloudConnection.dump_raw_catalog()`, which reads the configured catalog directly
 14from Airbyte. For documented self-managed deployments where the Public API root ends in
 15`/api/public/v1`, PyAirbyte infers the Config API root by replacing that suffix with
 16`/api/v1`.
 17
 18If your deployment uses custom ingress or a nonstandard reverse proxy, pass
 19`config_api_root` explicitly or set the `AIRBYTE_CLOUD_CONFIG_API_URL` environment
 20variable.
 21
 22```python
 23from airbyte import cloud
 24
 25workspace = cloud.CloudWorkspace(
 26    workspace_id="...",
 27    client_id="...",
 28    client_secret="...",
 29    api_root="https://airbyte.example.com/api/public/v1",
 30    config_api_root="https://airbyte.example.com/api/v1",
 31)
 32
 33connection = workspace.get_connection(connection_id="...")
 34raw_catalog = connection.dump_raw_catalog()
 35```
 36
 37## Examples
 38
 39### Basic Sync Example:
 40
 41```python
 42import airbyte as ab
 43from airbyte import cloud
 44
 45# Initialize an Airbyte Cloud workspace object
 46workspace = cloud.CloudWorkspace(
 47    workspace_id="123",
 48    api_key=ab.get_secret("AIRBYTE_CLOUD_API_KEY"),
 49)
 50
 51# Run a sync job on Airbyte Cloud
 52connection = workspace.get_connection(connection_id="456")
 53sync_result = connection.run_sync()
 54print(sync_result.get_job_status())
 55```
 56
 57### Example Read From Cloud Destination:
 58
 59If your destination is supported, you can read records directly from the
 60`SyncResult` object. Currently this is supported in Snowflake and BigQuery only.
 61
 62
 63```python
 64# Assuming we've already created a `connection` object...
 65
 66# Get the latest job result and print the stream names
 67sync_result = connection.get_sync_result()
 68print(sync_result.stream_names)
 69
 70# Get a dataset from the sync result
 71dataset: CachedDataset = sync_result.get_dataset("users")
 72
 73# Get a SQLAlchemy table to use in SQL queries...
 74users_table = dataset.to_sql_table()
 75print(f"Table name: {users_table.name}")
 76
 77# Or iterate over the dataset directly
 78for record in dataset:
 79    print(record)
 80```
 81"""
 82
 83from __future__ import annotations
 84
 85from typing import TYPE_CHECKING
 86
 87from airbyte.cloud.client import CloudClient
 88from airbyte.cloud.client_config import CloudClientConfig
 89from airbyte.cloud.connections import CloudConnection
 90from airbyte.cloud.models import (
 91    CloudDefaultContextInfo,
 92    CloudWorkspaceInfo,
 93    JobStatusEnum,
 94    JobTypeEnum,
 95    WorkspacePrivilegeScope,
 96)
 97from airbyte.cloud.organizations import CloudOrganization
 98from airbyte.cloud.sync_results import SyncResult
 99from airbyte.cloud.workspaces import CloudWorkspace
100
101
102# Submodules imported here for documentation reasons: https://github.com/mitmproxy/pdoc/issues/757
103if TYPE_CHECKING:
104    # ruff: noqa: TC004
105    from airbyte.cloud import (
106        client,
107        client_config,
108        connections,
109        constants,
110        organizations,
111        sync_results,
112        workspaces,
113    )
114
115
116__all__ = [
117    # Submodules
118    "workspaces",
119    "client",
120    "organizations",
121    "connections",
122    "constants",
123    "client_config",
124    "sync_results",
125    # Classes
126    "CloudClient",
127    "CloudOrganization",
128    "CloudWorkspace",
129    "CloudConnection",
130    "CloudClientConfig",
131    "CloudDefaultContextInfo",
132    "CloudWorkspaceInfo",
133    "SyncResult",
134    # Enums
135    "JobStatusEnum",
136    "JobTypeEnum",
137    "WorkspacePrivilegeScope",
138]
@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) -> 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) -> 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) -> 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: WorkspacePrivilegeScope = <WorkspacePrivilegeScope.MEMBER_OF: 'member_of'>, all_organizations: bool = False) -> list[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) -> 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[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) -> 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.

class CloudOrganization:
 22class CloudOrganization:
 23    """Information about an organization in Airbyte Cloud.
 24
 25    This class provides lazy loading of organization attributes including billing status.
 26    It is typically created via `CloudWorkspace.get_organization()`.
 27    """
 28
 29    def __init__(
 30        self,
 31        organization_id: str,
 32        organization_name: str | None = None,
 33        email: str | None = None,
 34        *,
 35        client_id: str | SecretString | None = None,
 36        client_secret: str | SecretString | None = None,
 37        bearer_token: str | SecretString | None = None,
 38        public_api_root: str | None = None,
 39        config_api_root: str | None = None,
 40    ) -> None:
 41        """Initialize a `CloudOrganization`."""
 42        self.organization_id = organization_id
 43        """The organization ID."""
 44
 45        self._organization_name = organization_name
 46        """Display name of the organization."""
 47
 48        self._email = email
 49        """Email associated with the organization."""
 50
 51        self._credentials = _AirbyteCredentials(
 52            client_id=SecretString(client_id) if client_id else None,
 53            client_secret=SecretString(client_secret) if client_secret else None,
 54            bearer_token=SecretString(bearer_token) if bearer_token else None,
 55            public_api_root=public_api_root or api_util.CLOUD_API_ROOT,
 56            config_api_root=config_api_root,
 57            organization_id=organization_id,
 58        )
 59        self._organization_info: dict[str, Any] | None = None
 60        self._organization_info_fetch_failed: bool = False
 61
 62    def _fetch_organization_info(self, *, force_refresh: bool = False) -> dict[str, Any]:
 63        """Fetch and cache organization info including billing status."""
 64        if force_refresh:
 65            self._organization_info_fetch_failed = False
 66
 67        if self._organization_info_fetch_failed and self._organization_info is None:
 68            return {}
 69
 70        if not force_refresh and self._organization_info is not None:
 71            return self._organization_info
 72
 73        try:
 74            self._organization_info = api_util.get_organization_info(
 75                organization_id=self.organization_id,
 76                api_root=self._credentials.public_api_root,
 77                config_api_root=self._credentials.config_api_root,
 78                client_id=self._credentials.client_id,
 79                client_secret=self._credentials.client_secret,
 80                bearer_token=self._credentials.bearer_token,
 81            )
 82        except Exception as ex:
 83            logger.debug("Failed to fetch organization info.", exc_info=ex)
 84            if self._organization_info is None:
 85                self._organization_info_fetch_failed = True
 86            return self._organization_info or {}
 87        else:
 88            return self._organization_info
 89
 90    @property
 91    def organization_name(self) -> str | None:
 92        """Display name of the organization."""
 93        if self._organization_name is not None:
 94            return self._organization_name
 95        info = self._fetch_organization_info()
 96        return info.get("organizationName")
 97
 98    @property
 99    def email(self) -> str | None:
100        """Email associated with the organization."""
101        if self._email is not None:
102            return self._email
103        info = self._fetch_organization_info()
104        return info.get("email")
105
106    def get_billing_status(self) -> CloudOrganizationBillingInfo:
107        """Fetch billing status for the organization or raise on failure."""
108        try:
109            info = api_util.get_organization_info(
110                organization_id=self.organization_id,
111                api_root=self._credentials.public_api_root,
112                config_api_root=self._credentials.config_api_root,
113                client_id=self._credentials.client_id,
114                client_secret=self._credentials.client_secret,
115                bearer_token=self._credentials.bearer_token,
116            )
117        except (requests.RequestException, ValueError) as ex:
118            raise AirbyteError(
119                message="Failed to retrieve organization billing information.",
120                context={"organization_id": self.organization_id},
121            ) from ex
122        billing = info.get("billing")
123        if not isinstance(billing, dict):
124            raise AirbyteError(
125                message="Organization info did not include billing details.",
126                context={"organization_id": self.organization_id},
127            )
128        payment_status = billing.get("paymentStatus")
129        subscription_status = billing.get("subscriptionStatus")
130        return CloudOrganizationBillingInfo(
131            payment_status=payment_status if isinstance(payment_status, str) else None,
132            subscription_status=(
133                subscription_status if isinstance(subscription_status, str) else None
134            ),
135            is_account_locked=api_util.is_account_locked(payment_status, subscription_status),
136        )
137
138    @property
139    def payment_status(self) -> str | None:
140        """Payment status of the organization."""
141        info = self._fetch_organization_info()
142        return (info.get("billing") or {}).get("paymentStatus")
143
144    @property
145    def subscription_status(self) -> str | None:
146        """Subscription status of the organization."""
147        info = self._fetch_organization_info()
148        return (info.get("billing") or {}).get("subscriptionStatus")
149
150    @property
151    def is_account_locked(self) -> bool:
152        """Whether the account is locked due to billing issues."""
153        return api_util.is_account_locked(self.payment_status, self.subscription_status)

Information about an organization in Airbyte Cloud.

This class provides lazy loading of organization attributes including billing status. It is typically created via CloudWorkspace.get_organization().

CloudOrganization( organization_id: str, organization_name: str | None = None, email: 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)
29    def __init__(
30        self,
31        organization_id: str,
32        organization_name: str | None = None,
33        email: str | None = None,
34        *,
35        client_id: str | SecretString | None = None,
36        client_secret: str | SecretString | None = None,
37        bearer_token: str | SecretString | None = None,
38        public_api_root: str | None = None,
39        config_api_root: str | None = None,
40    ) -> None:
41        """Initialize a `CloudOrganization`."""
42        self.organization_id = organization_id
43        """The organization ID."""
44
45        self._organization_name = organization_name
46        """Display name of the organization."""
47
48        self._email = email
49        """Email associated with the organization."""
50
51        self._credentials = _AirbyteCredentials(
52            client_id=SecretString(client_id) if client_id else None,
53            client_secret=SecretString(client_secret) if client_secret else None,
54            bearer_token=SecretString(bearer_token) if bearer_token else None,
55            public_api_root=public_api_root or api_util.CLOUD_API_ROOT,
56            config_api_root=config_api_root,
57            organization_id=organization_id,
58        )
59        self._organization_info: dict[str, Any] | None = None
60        self._organization_info_fetch_failed: bool = False

Initialize a CloudOrganization.

organization_id

The organization ID.

organization_name: str | None
90    @property
91    def organization_name(self) -> str | None:
92        """Display name of the organization."""
93        if self._organization_name is not None:
94            return self._organization_name
95        info = self._fetch_organization_info()
96        return info.get("organizationName")

Display name of the organization.

email: str | None
 98    @property
 99    def email(self) -> str | None:
100        """Email associated with the organization."""
101        if self._email is not None:
102            return self._email
103        info = self._fetch_organization_info()
104        return info.get("email")

Email associated with the organization.

def get_billing_status(self) -> airbyte.cloud.models.CloudOrganizationBillingInfo:
106    def get_billing_status(self) -> CloudOrganizationBillingInfo:
107        """Fetch billing status for the organization or raise on failure."""
108        try:
109            info = api_util.get_organization_info(
110                organization_id=self.organization_id,
111                api_root=self._credentials.public_api_root,
112                config_api_root=self._credentials.config_api_root,
113                client_id=self._credentials.client_id,
114                client_secret=self._credentials.client_secret,
115                bearer_token=self._credentials.bearer_token,
116            )
117        except (requests.RequestException, ValueError) as ex:
118            raise AirbyteError(
119                message="Failed to retrieve organization billing information.",
120                context={"organization_id": self.organization_id},
121            ) from ex
122        billing = info.get("billing")
123        if not isinstance(billing, dict):
124            raise AirbyteError(
125                message="Organization info did not include billing details.",
126                context={"organization_id": self.organization_id},
127            )
128        payment_status = billing.get("paymentStatus")
129        subscription_status = billing.get("subscriptionStatus")
130        return CloudOrganizationBillingInfo(
131            payment_status=payment_status if isinstance(payment_status, str) else None,
132            subscription_status=(
133                subscription_status if isinstance(subscription_status, str) else None
134            ),
135            is_account_locked=api_util.is_account_locked(payment_status, subscription_status),
136        )

Fetch billing status for the organization or raise on failure.

payment_status: str | None
138    @property
139    def payment_status(self) -> str | None:
140        """Payment status of the organization."""
141        info = self._fetch_organization_info()
142        return (info.get("billing") or {}).get("paymentStatus")

Payment status of the organization.

subscription_status: str | None
144    @property
145    def subscription_status(self) -> str | None:
146        """Subscription status of the organization."""
147        info = self._fetch_organization_info()
148        return (info.get("billing") or {}).get("subscriptionStatus")

Subscription status of the organization.

is_account_locked: bool
150    @property
151    def is_account_locked(self) -> bool:
152        """Whether the account is locked due to billing issues."""
153        return api_util.is_account_locked(self.payment_status, self.subscription_status)

Whether the account is locked due to billing issues.

@dataclass(init=False, kw_only=True)
class CloudWorkspace:
 111@dataclass(init=False, kw_only=True)  # noqa: PLR0904  # Core cloud API facade.
 112class CloudWorkspace:
 113    """A remote workspace on the Airbyte Cloud.
 114
 115    By overriding `api_root`, you can use this class to interact with self-managed Airbyte
 116    instances, both OSS and Enterprise.
 117
 118    Two authentication methods are supported (mutually exclusive):
 119    1. OAuth2 client credentials (client_id + client_secret)
 120    2. Bearer token authentication
 121
 122    Example with client credentials:
 123        ```python
 124        workspace = CloudWorkspace(
 125            workspace_id="...",
 126            client_id="...",
 127            client_secret="...",
 128        )
 129        ```
 130
 131    Example with bearer token:
 132        ```python
 133        workspace = CloudWorkspace(
 134            workspace_id="...",
 135            bearer_token="...",
 136        )
 137        ```
 138    """
 139
 140    workspace_id: str
 141    client_id: SecretString | None
 142    client_secret: SecretString | None
 143    api_root: str
 144    config_api_root: str | None
 145    """The Config API root URL."""
 146    bearer_token: SecretString | None
 147
 148    # Internal credentials objects (set in __init__, excluded from repr)
 149    _credentials: _AirbyteCredentials = field(init=False, repr=False)
 150    _client_config: CloudClientConfig = field(init=False, repr=False)
 151
 152    def __init__(
 153        self,
 154        *,
 155        workspace_id: str | None = None,
 156        client_id: str | SecretString | None = None,
 157        client_secret: str | SecretString | None = None,
 158        api_root: str | None = None,
 159        config_api_root: str | None = None,
 160        bearer_token: str | SecretString | None = None,
 161    ) -> None:
 162        """Validate and initialize credentials."""
 163        env_vars = not (client_id or client_secret or bearer_token)
 164        credentials = _AirbyteCredentials.from_auth(
 165            workspace_id=workspace_id,
 166            client_id=client_id,
 167            client_secret=client_secret,
 168            bearer_token=bearer_token,
 169            public_api_root=api_root,
 170            config_api_root=config_api_root,
 171            env_vars=env_vars,
 172        )
 173        if not credentials.workspace_id:
 174            raise exc.PyAirbyteInputError(
 175                message="Workspace ID is required.",
 176                guidance=(
 177                    "Provide a workspace ID, or call `get_default_cloud_context` to discover "
 178                    "available workspaces."
 179                ),
 180            )
 181
 182        self._credentials = credentials
 183        self.workspace_id = credentials.workspace_id or ""
 184        self.client_id = credentials.client_id
 185        self.client_secret = credentials.client_secret
 186        self.bearer_token = credentials.bearer_token
 187        self.api_root = credentials.public_api_root
 188        self.config_api_root = credentials.config_api_root
 189
 190        # Create internal CloudClientConfig object (validates mutual exclusivity)
 191        self._client_config = CloudClientConfig(
 192            client_id=self.client_id,
 193            client_secret=self.client_secret,
 194            bearer_token=self.bearer_token,
 195            api_root=self.api_root,
 196            config_api_root=self.config_api_root,
 197        )
 198
 199    @classmethod
 200    def from_env(
 201        cls,
 202        workspace_id: str | None = None,
 203        *,
 204        api_root: str | None = None,
 205        config_api_root: str | None = None,
 206    ) -> CloudWorkspace:
 207        """Create a CloudWorkspace using credentials from environment variables.
 208
 209        This factory method resolves credentials from environment variables,
 210        providing a convenient way to create a workspace without explicitly
 211        passing credentials.
 212
 213        Two authentication methods are supported (mutually exclusive):
 214        1. Bearer token (checked first)
 215        2. OAuth2 client credentials (fallback)
 216
 217        Environment variables used:
 218            - `AIRBYTE_CLOUD_BEARER_TOKEN`: Bearer token (alternative to client credentials).
 219            - `AIRBYTE_CLOUD_CLIENT_ID`: OAuth client ID (for client credentials flow).
 220            - `AIRBYTE_CLOUD_CLIENT_SECRET`: OAuth client secret (for client credentials flow).
 221            - `AIRBYTE_CLOUD_WORKSPACE_ID`: The workspace ID (if not passed as argument).
 222            - `AIRBYTE_CLOUD_API_URL`: Optional. The API root URL (defaults to Airbyte Cloud).
 223            - `AIRBYTE_CLOUD_CONFIG_API_URL`: Optional. The Config API root URL.
 224
 225        Args:
 226            workspace_id: The workspace ID. If not provided, will be resolved from
 227                the `AIRBYTE_CLOUD_WORKSPACE_ID` environment variable.
 228            api_root: The API root URL. If not provided, will be resolved from
 229                the `AIRBYTE_CLOUD_API_URL` environment variable, or default to
 230                the Airbyte Cloud API.
 231            config_api_root: The Config API root URL. If not provided, will be resolved
 232                from the `AIRBYTE_CLOUD_CONFIG_API_URL` environment variable.
 233
 234        Returns:
 235            A CloudWorkspace instance configured with credentials from the environment.
 236
 237        Raises:
 238            PyAirbyteInputError: If required credentials are not found in
 239                the environment or are incomplete.
 240
 241        Example:
 242            ```python
 243            # With workspace_id from environment
 244            workspace = CloudWorkspace.from_env()
 245
 246            # With explicit workspace_id
 247            workspace = CloudWorkspace.from_env(workspace_id="your-workspace-id")
 248            ```
 249        """
 250        return cls(
 251            workspace_id=workspace_id,
 252            api_root=api_root,
 253            config_api_root=config_api_root,
 254        )
 255
 256    @property
 257    def workspace_url(self) -> str | None:
 258        """The web URL of the workspace."""
 259        return f"{get_web_url_root(self.api_root)}/workspaces/{self.workspace_id}"
 260
 261    @cached_property
 262    def _organization_info(self) -> dict[str, Any]:
 263        """Fetch and cache organization info for this workspace.
 264
 265        Uses the Config API endpoint for an efficient O(1) lookup.
 266        This is an internal method; use get_organization() for public access.
 267        """
 268        return api_util.get_workspace_organization_info(
 269            workspace_id=self.workspace_id,
 270            api_root=self.api_root,
 271            config_api_root=self.config_api_root,
 272            client_id=self.client_id,
 273            client_secret=self.client_secret,
 274            bearer_token=self.bearer_token,
 275        )
 276
 277    @overload
 278    def get_organization(self) -> CloudOrganization: ...
 279
 280    @overload
 281    def get_organization(
 282        self,
 283        *,
 284        raise_on_error: Literal[True],
 285    ) -> CloudOrganization: ...
 286
 287    @overload
 288    def get_organization(
 289        self,
 290        *,
 291        raise_on_error: Literal[False],
 292    ) -> CloudOrganization | None: ...
 293
 294    def get_organization(
 295        self,
 296        *,
 297        raise_on_error: bool = True,
 298    ) -> CloudOrganization | None:
 299        """Get the organization this workspace belongs to.
 300
 301        Fetching organization info requires ORGANIZATION_READER permissions on the organization,
 302        which may not be available with workspace-scoped credentials.
 303
 304        Args:
 305            raise_on_error: If True (default), raises AirbyteError on permission or API errors.
 306                If False, returns None instead of raising.
 307
 308        Returns:
 309            CloudOrganization object with organization_id and organization_name,
 310            or None if raise_on_error=False and an error occurred.
 311
 312        Raises:
 313            AirbyteError: If raise_on_error=True and the organization info cannot be fetched
 314                (e.g., due to insufficient permissions or missing data).
 315        """
 316        try:
 317            info = self._organization_info
 318        except (AirbyteError, NotImplementedError):
 319            if raise_on_error:
 320                raise
 321            return None
 322
 323        organization_id = info.get("organizationId")
 324        organization_name = info.get("organizationName")
 325
 326        # Validate that both organization_id and organization_name are non-null and non-empty
 327        if not organization_id or not organization_name:
 328            if raise_on_error:
 329                raise AirbyteError(
 330                    message="Organization info is incomplete.",
 331                    context={
 332                        "organization_id": organization_id,
 333                        "organization_name": organization_name,
 334                    },
 335                )
 336            return None
 337
 338        organization_credentials = self._credentials.with_organization_id(organization_id)
 339        return CloudOrganization(
 340            organization_id=organization_id,
 341            organization_name=organization_name,
 342            client_id=organization_credentials.client_id,
 343            client_secret=organization_credentials.client_secret,
 344            bearer_token=organization_credentials.bearer_token,
 345            public_api_root=organization_credentials.public_api_root,
 346            config_api_root=organization_credentials.config_api_root,
 347        )
 348
 349    # Test connection and creds
 350
 351    def connect(self) -> None:
 352        """Check that the workspace is reachable and raise an exception otherwise.
 353
 354        Note: It is not necessary to call this method before calling other operations. It
 355              serves primarily as a simple check to ensure that the workspace is reachable
 356              and credentials are correct.
 357        """
 358        _ = api_util.get_workspace(
 359            api_root=self.api_root,
 360            workspace_id=self.workspace_id,
 361            client_id=self.client_id,
 362            client_secret=self.client_secret,
 363            bearer_token=self.bearer_token,
 364        )
 365        print(f"Successfully connected to workspace: {self.workspace_url}")
 366
 367    # Get sources, destinations, and connections
 368
 369    def get_connection(
 370        self,
 371        connection_id: str,
 372    ) -> CloudConnection:
 373        """Get a connection by ID.
 374
 375        This method does not fetch data from the API. It returns a `CloudConnection` object,
 376        which will be loaded lazily as needed.
 377        """
 378        return CloudConnection(
 379            workspace=self,
 380            connection_id=connection_id,
 381        )
 382
 383    def get_source(
 384        self,
 385        source_id: str,
 386    ) -> CloudSource:
 387        """Get a source by ID.
 388
 389        This method does not fetch data from the API. It returns a `CloudSource` object,
 390        which will be loaded lazily as needed.
 391        """
 392        return CloudSource(
 393            workspace=self,
 394            connector_id=source_id,
 395        )
 396
 397    def get_destination(
 398        self,
 399        destination_id: str,
 400    ) -> CloudDestination:
 401        """Get a destination by ID.
 402
 403        This method does not fetch data from the API. It returns a `CloudDestination` object,
 404        which will be loaded lazily as needed.
 405        """
 406        return CloudDestination(
 407            workspace=self,
 408            connector_id=destination_id,
 409        )
 410
 411    def check_connector_setup(
 412        self,
 413        connector_type: Literal["source", "destination"],
 414        connector_id: str,
 415    ) -> CheckResult:
 416        """Run one connection check on a connector that belongs to this workspace.
 417
 418        Confirms a person has finished a deferred-credential setup in Airbyte Cloud. The
 419        connector's workspace is verified first so a check can never be run against a connector
 420        outside this workspace.
 421        """
 422        connector: CloudSource | CloudDestination
 423        if connector_type == "source":
 424            owner_id = api_util.get_source(
 425                source_id=connector_id,
 426                api_root=self.api_root,
 427                client_id=self.client_id,
 428                client_secret=self.client_secret,
 429                bearer_token=self.bearer_token,
 430            ).workspace_id
 431            connector = self.get_source(connector_id)
 432        else:
 433            owner_id = api_util.get_destination(
 434                destination_id=connector_id,
 435                api_root=self.api_root,
 436                client_id=self.client_id,
 437                client_secret=self.client_secret,
 438                bearer_token=self.bearer_token,
 439            ).workspace_id
 440            connector = self.get_destination(connector_id)
 441        if owner_id != self.workspace_id:
 442            raise exc.AirbyteMissingResourceError(
 443                resource_type=connector_type,
 444                resource_name_or_id=connector_id,
 445                context={"workspace_id": self.workspace_id},
 446            )
 447        try:
 448            return connector.check(raise_on_error=False)
 449        except AirbyteError as ex:
 450            status_code = (ex.context or {}).get("status_code")
 451            if status_code == HTTPStatus.UNPROCESSABLE_ENTITY:
 452                return CheckResult(success=False)
 453            raise AirbyteError(
 454                message="Cloud could not check the connector setup.",
 455                context={"status_code": status_code},
 456            ) from None
 457
 458    # Deploy sources and destinations
 459
 460    def deploy_source(
 461        self,
 462        name: str,
 463        source: Source | dict[str, Any],
 464        *,
 465        unique: bool = True,
 466        random_name_suffix: bool = False,
 467        definition_id: str | None = None,
 468        defer_credentials: bool = False,
 469    ) -> CloudSource:
 470        """Deploy a source to the workspace.
 471
 472        Returns the newly deployed source.
 473
 474        Args:
 475            name: The name to use when deploying.
 476            source: The source object to deploy, or (with `defer_credentials=True`) a
 477                dictionary of non-secret configuration values.
 478            unique: Whether to require a unique name. If `True`, duplicate names
 479                are not allowed. Defaults to `True`.
 480            random_name_suffix: Whether to append a random suffix to the name.
 481            definition_id: The source definition ID. Required with `defer_credentials=True`.
 482            defer_credentials: Save a draft with partial configuration. A person completes
 483                credentials and other missing settings at the returned source's `connector_url`.
 484                A successful connection check promotes the draft. Raises
 485                `AirbyteDeferredSetupError` if Cloud does not acknowledge draft mode.
 486        """
 487        if defer_credentials:
 488            return CloudSource(
 489                workspace=self,
 490                connector_id=self._deploy_deferred(
 491                    connector_type="source",
 492                    name=name,
 493                    config=source,
 494                    definition_id=definition_id,
 495                    unique=unique,
 496                    random_name_suffix=random_name_suffix,
 497                ),
 498            )
 499        if isinstance(source, dict):
 500            raise exc.PyAirbyteInputError(
 501                message="`source` must be a `Source` object unless `defer_credentials=True`.",
 502            )
 503
 504        source_config_dict = source._hydrated_config.copy()  # noqa: SLF001 (non-public API)
 505        source_config_dict["sourceType"] = source.name.replace("source-", "")
 506
 507        if random_name_suffix:
 508            name += f" (ID: {text_util.generate_random_suffix()})"
 509
 510        if unique:
 511            existing = self.list_sources(name=name)
 512            if existing:
 513                raise exc.AirbyteDuplicateResourcesError(
 514                    resource_type="source",
 515                    resource_name=name,
 516                )
 517
 518        deployed_source = api_util.create_source(
 519            name=name,
 520            api_root=self.api_root,
 521            workspace_id=self.workspace_id,
 522            config=source_config_dict,
 523            definition_id=definition_id,
 524            client_id=self.client_id,
 525            client_secret=self.client_secret,
 526            bearer_token=self.bearer_token,
 527        )
 528        return CloudSource(
 529            workspace=self,
 530            connector_id=deployed_source.source_id,
 531        )
 532
 533    def deploy_destination(
 534        self,
 535        name: str,
 536        destination: Destination | dict[str, Any],
 537        *,
 538        unique: bool = True,
 539        random_name_suffix: bool = False,
 540        definition_id: str | None = None,
 541        defer_credentials: bool = False,
 542    ) -> CloudDestination:
 543        """Deploy a destination to the workspace.
 544
 545        Returns the newly deployed destination ID.
 546
 547        Args:
 548            name: The name to use when deploying.
 549            destination: The destination to deploy. Can be a local Airbyte `Destination` object or a
 550                dictionary of configuration values.
 551            unique: Whether to require a unique name. If `True`, duplicate names
 552                are not allowed. Defaults to `True`.
 553            random_name_suffix: Whether to append a random suffix to the name.
 554            definition_id: The destination definition ID. Required with `defer_credentials=True`;
 555                otherwise the type is inferred from `destinationType`.
 556            defer_credentials: Create the destination without its credentials. See
 557                `deploy_source`.
 558        """
 559        if defer_credentials:
 560            return CloudDestination(
 561                workspace=self,
 562                connector_id=self._deploy_deferred(
 563                    connector_type="destination",
 564                    name=name,
 565                    config=destination,
 566                    definition_id=definition_id,
 567                    unique=unique,
 568                    random_name_suffix=random_name_suffix,
 569                ),
 570            )
 571
 572        if isinstance(destination, Destination):
 573            destination_conf_dict = destination._hydrated_config.copy()  # noqa: SLF001 (non-public API)
 574            destination_conf_dict["destinationType"] = destination.name.replace("destination-", "")
 575            # raise ValueError(destination_conf_dict)
 576        else:
 577            destination_conf_dict = destination.copy()
 578            if "destinationType" not in destination_conf_dict:
 579                raise exc.PyAirbyteInputError(
 580                    message="Missing `destinationType` in configuration dictionary.",
 581                )
 582
 583        if random_name_suffix:
 584            name += f" (ID: {text_util.generate_random_suffix()})"
 585
 586        if unique:
 587            existing = self.list_destinations(name=name)
 588            if existing:
 589                raise exc.AirbyteDuplicateResourcesError(
 590                    resource_type="destination",
 591                    resource_name=name,
 592                )
 593
 594        deployed_destination = api_util.create_destination(
 595            name=name,
 596            api_root=self.api_root,
 597            workspace_id=self.workspace_id,
 598            config=destination_conf_dict,  # Wants a dataclass but accepts dict
 599            client_id=self.client_id,
 600            client_secret=self.client_secret,
 601            bearer_token=self.bearer_token,
 602        )
 603        return CloudDestination(
 604            workspace=self,
 605            connector_id=deployed_destination.destination_id,
 606        )
 607
 608    def _deploy_deferred(
 609        self,
 610        *,
 611        connector_type: Literal["source", "destination"],
 612        name: str,
 613        config: object,
 614        definition_id: str | None,
 615        unique: bool,
 616        random_name_suffix: bool,
 617    ) -> str:
 618        """Create a connector with deferred credentials on the Config API and return its ID."""
 619        config_dict, definition_id = _deferred_credentials_config(
 620            config, definition_id=definition_id
 621        )
 622
 623        if random_name_suffix:
 624            name += f" (ID: {text_util.generate_random_suffix()})"
 625
 626        if unique:
 627            existing = (
 628                self.list_sources(name=name)
 629                if connector_type == "source"
 630                else self.list_destinations(name=name)
 631            )
 632            if existing:
 633                raise exc.AirbyteDuplicateResourcesError(
 634                    resource_type=connector_type,
 635                    resource_name=name,
 636                )
 637
 638        return api_util.create_connector_deferred(
 639            connector_type=connector_type,
 640            name=name,
 641            workspace_id=self.workspace_id,
 642            definition_id=definition_id,
 643            config=config_dict,
 644            api_root=self.api_root,
 645            config_api_root=self.config_api_root,
 646            client_id=self.client_id,
 647            client_secret=self.client_secret,
 648            bearer_token=self.bearer_token,
 649        )
 650
 651    def permanently_delete_source(
 652        self,
 653        source: str | CloudSource,
 654        *,
 655        safe_mode: bool = True,
 656    ) -> None:
 657        """Delete a source from the workspace.
 658
 659        You can pass either the source ID `str` or a deployed `Source` object.
 660
 661        Args:
 662            source: The source ID or CloudSource object to delete
 663            safe_mode: If True, requires the source name to contain "delete-me" or "deleteme"
 664                (case insensitive) to prevent accidental deletion. Defaults to True.
 665        """
 666        if not isinstance(source, (str, CloudSource)):
 667            raise exc.PyAirbyteInputError(
 668                message="Invalid source type.",
 669                input_value=type(source).__name__,
 670            )
 671
 672        api_util.delete_source(
 673            source_id=source.connector_id if isinstance(source, CloudSource) else source,
 674            source_name=source.name if isinstance(source, CloudSource) else None,
 675            api_root=self.api_root,
 676            client_id=self.client_id,
 677            client_secret=self.client_secret,
 678            bearer_token=self.bearer_token,
 679            safe_mode=safe_mode,
 680        )
 681
 682    # Deploy and delete destinations
 683
 684    def permanently_delete_destination(
 685        self,
 686        destination: str | CloudDestination,
 687        *,
 688        safe_mode: bool = True,
 689    ) -> None:
 690        """Delete a deployed destination from the workspace.
 691
 692        You can pass either the `Cache` class or the deployed destination ID as a `str`.
 693
 694        Args:
 695            destination: The destination ID or CloudDestination object to delete
 696            safe_mode: If True, requires the destination name to contain "delete-me" or "deleteme"
 697                (case insensitive) to prevent accidental deletion. Defaults to True.
 698        """
 699        if not isinstance(destination, (str, CloudDestination)):
 700            raise exc.PyAirbyteInputError(
 701                message="Invalid destination type.",
 702                input_value=type(destination).__name__,
 703            )
 704
 705        api_util.delete_destination(
 706            destination_id=(
 707                destination if isinstance(destination, str) else destination.destination_id
 708            ),
 709            destination_name=(
 710                destination.name if isinstance(destination, CloudDestination) else None
 711            ),
 712            api_root=self.api_root,
 713            client_id=self.client_id,
 714            client_secret=self.client_secret,
 715            bearer_token=self.bearer_token,
 716            safe_mode=safe_mode,
 717        )
 718
 719    # Deploy and delete connections
 720
 721    def deploy_connection(
 722        self,
 723        connection_name: str,
 724        *,
 725        source: CloudSource | str,
 726        selected_streams: list[str],
 727        destination: CloudDestination | str,
 728        table_prefix: str | None = None,
 729    ) -> CloudConnection:
 730        """Create a new connection between an already deployed source and destination.
 731
 732        Returns the newly deployed connection object.
 733
 734        Args:
 735            connection_name: The name of the connection.
 736            source: The deployed source. You can pass a source ID or a CloudSource object.
 737            destination: The deployed destination. You can pass a destination ID or a
 738                CloudDestination object.
 739            table_prefix: Optional. The table prefix to use when syncing to the destination.
 740            selected_streams: The selected stream names to sync within the connection.
 741        """
 742        if not selected_streams:
 743            raise exc.PyAirbyteInputError(
 744                guidance="You must provide `selected_streams` when creating a connection."
 745            )
 746
 747        source_id: str = source if isinstance(source, str) else source.connector_id
 748        destination_id: str = (
 749            destination if isinstance(destination, str) else destination.connector_id
 750        )
 751
 752        deployed_connection = api_util.create_connection(
 753            name=connection_name,
 754            source_id=source_id,
 755            destination_id=destination_id,
 756            api_root=self.api_root,
 757            workspace_id=self.workspace_id,
 758            selected_stream_names=selected_streams,
 759            prefix=table_prefix or "",
 760            client_id=self.client_id,
 761            client_secret=self.client_secret,
 762            bearer_token=self.bearer_token,
 763        )
 764
 765        return CloudConnection(
 766            workspace=self,
 767            connection_id=deployed_connection.connection_id,
 768            source=deployed_connection.source_id,
 769            destination=deployed_connection.destination_id,
 770        )
 771
 772    def permanently_delete_connection(
 773        self,
 774        connection: str | CloudConnection,
 775        *,
 776        cascade_delete_source: bool = False,
 777        cascade_delete_destination: bool = False,
 778        safe_mode: bool = True,
 779    ) -> None:
 780        """Delete a deployed connection from the workspace.
 781
 782        Args:
 783            connection: The connection ID or CloudConnection object to delete
 784            cascade_delete_source: If True, also delete the source after deleting the connection
 785            cascade_delete_destination: If True, also delete the destination after deleting
 786                the connection
 787            safe_mode: If True, requires the connection name to contain "delete-me" or "deleteme"
 788                (case insensitive) to prevent accidental deletion. Defaults to True. Also applies
 789                to cascade deletes.
 790        """
 791        if connection is None:
 792            raise ValueError("No connection ID provided.")
 793
 794        if isinstance(connection, str):
 795            connection = CloudConnection(
 796                workspace=self,
 797                connection_id=connection,
 798            )
 799
 800        api_util.delete_connection(
 801            connection_id=connection.connection_id,
 802            connection_name=connection.name,
 803            api_root=self.api_root,
 804            workspace_id=self.workspace_id,
 805            client_id=self.client_id,
 806            client_secret=self.client_secret,
 807            bearer_token=self.bearer_token,
 808            safe_mode=safe_mode,
 809        )
 810
 811        if cascade_delete_source:
 812            self.permanently_delete_source(
 813                source=connection.source_id,
 814                safe_mode=safe_mode,
 815            )
 816        if cascade_delete_destination:
 817            self.permanently_delete_destination(
 818                destination=connection.destination_id,
 819                safe_mode=safe_mode,
 820            )
 821
 822    # List workspaces, sources, destinations, and connections
 823
 824    def list_workspaces(
 825        self,
 826        name: str | None = None,
 827        *,
 828        name_filter: Callable | None = None,
 829        limit: int | None = None,
 830    ) -> list[CloudWorkspaceInfo]:
 831        """List workspaces available to the current credentials, with an optional limit."""
 832        return [
 833            CloudWorkspaceInfo.from_api_response(workspace)
 834            for workspace in api_util.list_workspaces(
 835                workspace_id="",
 836                api_root=self.api_root,
 837                name=name,
 838                name_filter=name_filter,
 839                client_id=self.client_id,
 840                client_secret=self.client_secret,
 841                bearer_token=self.bearer_token,
 842                limit=limit,
 843            )
 844        ]
 845
 846    def rename(
 847        self,
 848        name: str,
 849    ) -> CloudWorkspace:
 850        """Rename this workspace."""
 851        api_util.rename_workspace(
 852            workspace_id=self.workspace_id,
 853            name=name,
 854            api_root=self.api_root,
 855            client_id=self.client_id,
 856            client_secret=self.client_secret,
 857            bearer_token=self.bearer_token,
 858        )
 859        return self
 860
 861    def permanently_delete(
 862        self,
 863        *,
 864        workspace_name: str | None = None,
 865        safe_mode: bool = True,
 866    ) -> None:
 867        """Permanently delete this workspace if it has no connections.
 868
 869        When `safe_mode` is enabled, the workspace name must contain `delete-me`
 870        or `deleteme`. This also checks for existing connections before deleting
 871        and raises `AirbyteWorkspaceNotEmptyError` if the workspace is not empty.
 872        """
 873        api_util.permanently_delete_workspace(
 874            workspace_id=self.workspace_id,
 875            workspace_name=workspace_name,
 876            api_root=self.api_root,
 877            client_id=self.client_id,
 878            client_secret=self.client_secret,
 879            bearer_token=self.bearer_token,
 880            safe_mode=safe_mode,
 881        )
 882
 883    def list_connections(
 884        self,
 885        name: str | None = None,
 886        *,
 887        name_filter: Callable | None = None,
 888        limit: int | None = None,
 889    ) -> list[CloudConnection]:
 890        """List connections by name in the workspace, with an optional limit."""
 891        connections = api_util.list_connections(
 892            api_root=self.api_root,
 893            workspace_id=self.workspace_id,
 894            name=name,
 895            name_filter=name_filter,
 896            limit=limit,
 897            client_id=self.client_id,
 898            client_secret=self.client_secret,
 899            bearer_token=self.bearer_token,
 900        )
 901        return [
 902            CloudConnection._from_connection_response(  # noqa: SLF001 (non-public API)
 903                workspace=self,
 904                connection_response=connection,
 905            )
 906            for connection in connections
 907        ]
 908
 909    def list_sources(
 910        self,
 911        name: str | None = None,
 912        *,
 913        name_filter: Callable | None = None,
 914        limit: int | None = None,
 915    ) -> list[CloudSource]:
 916        """List all sources in the workspace, with an optional limit."""
 917        sources = api_util.list_sources(
 918            api_root=self.api_root,
 919            workspace_id=self.workspace_id,
 920            name=name,
 921            name_filter=name_filter,
 922            limit=limit,
 923            client_id=self.client_id,
 924            client_secret=self.client_secret,
 925            bearer_token=self.bearer_token,
 926        )
 927        return [
 928            CloudSource._from_source_response(  # noqa: SLF001 (non-public API)
 929                workspace=self,
 930                source_response=source,
 931            )
 932            for source in sources
 933        ]
 934
 935    def list_destinations(
 936        self,
 937        name: str | None = None,
 938        *,
 939        name_filter: Callable | None = None,
 940        limit: int | None = None,
 941    ) -> list[CloudDestination]:
 942        """List all destinations in the workspace, with an optional limit."""
 943        destinations = api_util.list_destinations(
 944            api_root=self.api_root,
 945            workspace_id=self.workspace_id,
 946            name=name,
 947            name_filter=name_filter,
 948            limit=limit,
 949            client_id=self.client_id,
 950            client_secret=self.client_secret,
 951            bearer_token=self.bearer_token,
 952        )
 953        return [
 954            CloudDestination._from_destination_response(  # noqa: SLF001 (non-public API)
 955                workspace=self,
 956                destination_response=destination,
 957            )
 958            for destination in destinations
 959        ]
 960
 961    def publish_custom_source_definition(
 962        self,
 963        name: str,
 964        *,
 965        manifest_yaml: dict[str, Any] | Path | str | None = None,
 966        docker_image: str | None = None,
 967        docker_tag: str | None = None,
 968        unique: bool = True,
 969        pre_validate: bool = True,
 970        testing_values: dict[str, Any] | None = None,
 971    ) -> CustomCloudSourceDefinition:
 972        """Publish a custom source connector definition.
 973
 974        You must specify EITHER manifest_yaml (for YAML connectors) OR both docker_image
 975        and docker_tag (for Docker connectors), but not both.
 976
 977        Args:
 978            name: Display name for the connector definition
 979            manifest_yaml: Low-code CDK manifest (dict, Path to YAML file, or YAML string)
 980            docker_image: Docker repository (e.g., 'airbyte/source-custom')
 981            docker_tag: Docker image tag (e.g., '1.0.0')
 982            unique: Whether to enforce name uniqueness
 983            pre_validate: Whether to validate manifest client-side (YAML only)
 984            testing_values: Optional configuration values to use for testing in the
 985                Connector Builder UI. If provided, these values are stored as the complete
 986                testing values object for the connector builder project (replaces any existing
 987                values), allowing immediate test read operations.
 988
 989        Returns:
 990            CustomCloudSourceDefinition object representing the created definition
 991
 992        Raises:
 993            PyAirbyteInputError: If both or neither of manifest_yaml and docker_image provided
 994            AirbyteDuplicateResourcesError: If unique=True and name already exists
 995        """
 996        is_yaml = manifest_yaml is not None
 997        is_docker = docker_image is not None
 998
 999        if is_yaml == is_docker:
1000            raise exc.PyAirbyteInputError(
1001                message=(
1002                    "Must specify EITHER manifest_yaml (for YAML connectors) OR "
1003                    "docker_image + docker_tag (for Docker connectors), but not both"
1004                ),
1005                context={
1006                    "manifest_yaml_provided": is_yaml,
1007                    "docker_image_provided": is_docker,
1008                },
1009            )
1010
1011        if is_docker and docker_tag is None:
1012            raise exc.PyAirbyteInputError(
1013                message="docker_tag is required when docker_image is specified",
1014                context={"docker_image": docker_image},
1015            )
1016
1017        if unique:
1018            existing = self.list_custom_source_definitions(
1019                definition_type="yaml" if is_yaml else "docker",
1020            )
1021            if any(d.name == name for d in existing):
1022                raise exc.AirbyteDuplicateResourcesError(
1023                    resource_type="custom_source_definition",
1024                    resource_name=name,
1025                )
1026
1027        if is_yaml:
1028            manifest_dict: dict[str, Any]
1029            if isinstance(manifest_yaml, Path):
1030                manifest_dict = yaml.safe_load(manifest_yaml.read_text())
1031            elif isinstance(manifest_yaml, str):
1032                manifest_dict = yaml.safe_load(manifest_yaml)
1033            elif manifest_yaml is not None:
1034                manifest_dict = manifest_yaml
1035            else:
1036                raise exc.PyAirbyteInputError(
1037                    message="manifest_yaml is required for YAML connectors",
1038                    context={"name": name},
1039                )
1040
1041            if pre_validate:
1042                api_util.validate_yaml_manifest(manifest_dict, raise_on_error=True)
1043
1044            result = api_util.create_custom_yaml_source_definition(
1045                name=name,
1046                workspace_id=self.workspace_id,
1047                manifest=manifest_dict,
1048                api_root=self.api_root,
1049                client_id=self.client_id,
1050                client_secret=self.client_secret,
1051                bearer_token=self.bearer_token,
1052            )
1053            custom_definition = CustomCloudSourceDefinition._from_yaml_response(  # noqa: SLF001
1054                self, result
1055            )
1056
1057            # Set testing values if provided
1058            if testing_values is not None:
1059                custom_definition.set_testing_values(testing_values)
1060
1061            return custom_definition
1062
1063        raise NotImplementedError(
1064            "Docker custom source definitions are not yet supported. "
1065            "Only YAML manifest-based custom sources are currently available."
1066        )
1067
1068    def list_custom_source_definitions(
1069        self,
1070        *,
1071        definition_type: Literal["yaml", "docker"],
1072    ) -> list[CustomCloudSourceDefinition]:
1073        """List custom source connector definitions.
1074
1075        Args:
1076            definition_type: Connector type to list ("yaml" or "docker"). Required.
1077
1078        Returns:
1079            List of CustomCloudSourceDefinition objects matching the specified type
1080        """
1081        if definition_type == "yaml":
1082            yaml_definitions = api_util.list_custom_yaml_source_definitions(
1083                workspace_id=self.workspace_id,
1084                api_root=self.api_root,
1085                client_id=self.client_id,
1086                client_secret=self.client_secret,
1087                bearer_token=self.bearer_token,
1088            )
1089            return [
1090                CustomCloudSourceDefinition._from_yaml_response(self, d)  # noqa: SLF001
1091                for d in yaml_definitions
1092            ]
1093
1094        raise NotImplementedError(
1095            "Docker custom source definitions are not yet supported. "
1096            "Only YAML manifest-based custom sources are currently available."
1097        )
1098
1099    def get_custom_source_definition(
1100        self,
1101        definition_id: str,
1102        *,
1103        definition_type: Literal["yaml", "docker"],
1104    ) -> CustomCloudSourceDefinition:
1105        """Get a specific custom source definition by ID.
1106
1107        Args:
1108            definition_id: The definition ID
1109            definition_type: Connector type ("yaml" or "docker"). Required.
1110
1111        Returns:
1112            CustomCloudSourceDefinition object
1113        """
1114        if definition_type == "yaml":
1115            result = api_util.get_custom_yaml_source_definition(
1116                workspace_id=self.workspace_id,
1117                definition_id=definition_id,
1118                api_root=self.api_root,
1119                client_id=self.client_id,
1120                client_secret=self.client_secret,
1121                bearer_token=self.bearer_token,
1122            )
1123            return CustomCloudSourceDefinition._from_yaml_response(self, result)  # noqa: SLF001
1124
1125        raise NotImplementedError(
1126            "Docker custom source definitions are not yet supported. "
1127            "Only YAML manifest-based custom sources are currently available."
1128        )

A remote workspace on the Airbyte Cloud.

By overriding api_root, you can use this class to interact with self-managed Airbyte instances, both OSS and Enterprise.

Two authentication methods are supported (mutually exclusive):

  1. OAuth2 client credentials (client_id + client_secret)
  2. Bearer token authentication
Example with client credentials:
workspace = CloudWorkspace(
    workspace_id="...",
    client_id="...",
    client_secret="...",
)
Example with bearer token:
workspace = CloudWorkspace(
    workspace_id="...",
    bearer_token="...",
)
CloudWorkspace( *, workspace_id: str | None = None, client_id: str | airbyte.secrets.SecretString | None = None, client_secret: str | airbyte.secrets.SecretString | None = None, api_root: str | None = None, config_api_root: str | None = None, bearer_token: str | airbyte.secrets.SecretString | None = None)
152    def __init__(
153        self,
154        *,
155        workspace_id: str | None = None,
156        client_id: str | SecretString | None = None,
157        client_secret: str | SecretString | None = None,
158        api_root: str | None = None,
159        config_api_root: str | None = None,
160        bearer_token: str | SecretString | None = None,
161    ) -> None:
162        """Validate and initialize credentials."""
163        env_vars = not (client_id or client_secret or bearer_token)
164        credentials = _AirbyteCredentials.from_auth(
165            workspace_id=workspace_id,
166            client_id=client_id,
167            client_secret=client_secret,
168            bearer_token=bearer_token,
169            public_api_root=api_root,
170            config_api_root=config_api_root,
171            env_vars=env_vars,
172        )
173        if not credentials.workspace_id:
174            raise exc.PyAirbyteInputError(
175                message="Workspace ID is required.",
176                guidance=(
177                    "Provide a workspace ID, or call `get_default_cloud_context` to discover "
178                    "available workspaces."
179                ),
180            )
181
182        self._credentials = credentials
183        self.workspace_id = credentials.workspace_id or ""
184        self.client_id = credentials.client_id
185        self.client_secret = credentials.client_secret
186        self.bearer_token = credentials.bearer_token
187        self.api_root = credentials.public_api_root
188        self.config_api_root = credentials.config_api_root
189
190        # Create internal CloudClientConfig object (validates mutual exclusivity)
191        self._client_config = CloudClientConfig(
192            client_id=self.client_id,
193            client_secret=self.client_secret,
194            bearer_token=self.bearer_token,
195            api_root=self.api_root,
196            config_api_root=self.config_api_root,
197        )

Validate and initialize credentials.

workspace_id: str
client_id: airbyte.secrets.SecretString | None
client_secret: airbyte.secrets.SecretString | None
api_root: str
config_api_root: str | None

The Config API root URL.

bearer_token: airbyte.secrets.SecretString | None
@classmethod
def from_env( cls, workspace_id: str | None = None, *, api_root: str | None = None, config_api_root: str | None = None) -> CloudWorkspace:
199    @classmethod
200    def from_env(
201        cls,
202        workspace_id: str | None = None,
203        *,
204        api_root: str | None = None,
205        config_api_root: str | None = None,
206    ) -> CloudWorkspace:
207        """Create a CloudWorkspace using credentials from environment variables.
208
209        This factory method resolves credentials from environment variables,
210        providing a convenient way to create a workspace without explicitly
211        passing credentials.
212
213        Two authentication methods are supported (mutually exclusive):
214        1. Bearer token (checked first)
215        2. OAuth2 client credentials (fallback)
216
217        Environment variables used:
218            - `AIRBYTE_CLOUD_BEARER_TOKEN`: Bearer token (alternative to client credentials).
219            - `AIRBYTE_CLOUD_CLIENT_ID`: OAuth client ID (for client credentials flow).
220            - `AIRBYTE_CLOUD_CLIENT_SECRET`: OAuth client secret (for client credentials flow).
221            - `AIRBYTE_CLOUD_WORKSPACE_ID`: The workspace ID (if not passed as argument).
222            - `AIRBYTE_CLOUD_API_URL`: Optional. The API root URL (defaults to Airbyte Cloud).
223            - `AIRBYTE_CLOUD_CONFIG_API_URL`: Optional. The Config API root URL.
224
225        Args:
226            workspace_id: The workspace ID. If not provided, will be resolved from
227                the `AIRBYTE_CLOUD_WORKSPACE_ID` environment variable.
228            api_root: The API root URL. If not provided, will be resolved from
229                the `AIRBYTE_CLOUD_API_URL` environment variable, or default to
230                the Airbyte Cloud API.
231            config_api_root: The Config API root URL. If not provided, will be resolved
232                from the `AIRBYTE_CLOUD_CONFIG_API_URL` environment variable.
233
234        Returns:
235            A CloudWorkspace instance configured with credentials from the environment.
236
237        Raises:
238            PyAirbyteInputError: If required credentials are not found in
239                the environment or are incomplete.
240
241        Example:
242            ```python
243            # With workspace_id from environment
244            workspace = CloudWorkspace.from_env()
245
246            # With explicit workspace_id
247            workspace = CloudWorkspace.from_env(workspace_id="your-workspace-id")
248            ```
249        """
250        return cls(
251            workspace_id=workspace_id,
252            api_root=api_root,
253            config_api_root=config_api_root,
254        )

Create a CloudWorkspace using credentials from environment variables.

This factory method resolves credentials from environment variables, providing a convenient way to create a workspace without explicitly passing credentials.

Two authentication methods are supported (mutually exclusive):

  1. Bearer token (checked first)
  2. OAuth2 client credentials (fallback)
Environment variables used:
  • AIRBYTE_CLOUD_BEARER_TOKEN: Bearer token (alternative to client credentials).
  • AIRBYTE_CLOUD_CLIENT_ID: OAuth client ID (for client credentials flow).
  • AIRBYTE_CLOUD_CLIENT_SECRET: OAuth client secret (for client credentials flow).
  • AIRBYTE_CLOUD_WORKSPACE_ID: The workspace ID (if not passed as argument).
  • AIRBYTE_CLOUD_API_URL: Optional. The API root URL (defaults to Airbyte Cloud).
  • AIRBYTE_CLOUD_CONFIG_API_URL: Optional. The Config API root URL.
Arguments:
  • workspace_id: The workspace ID. If not provided, will be resolved from the AIRBYTE_CLOUD_WORKSPACE_ID environment variable.
  • api_root: The API root URL. If not provided, will be resolved from the AIRBYTE_CLOUD_API_URL environment variable, or default to the Airbyte Cloud API.
  • config_api_root: The Config API root URL. If not provided, will be resolved from the AIRBYTE_CLOUD_CONFIG_API_URL environment variable.
Returns:

A CloudWorkspace instance configured with credentials from the environment.

Raises:
  • PyAirbyteInputError: If required credentials are not found in the environment or are incomplete.
Example:
# With workspace_id from environment
workspace = CloudWorkspace.from_env()

# With explicit workspace_id
workspace = CloudWorkspace.from_env(workspace_id="your-workspace-id")
workspace_url: str | None
256    @property
257    def workspace_url(self) -> str | None:
258        """The web URL of the workspace."""
259        return f"{get_web_url_root(self.api_root)}/workspaces/{self.workspace_id}"

The web URL of the workspace.

def get_organization( self, *, raise_on_error: bool = True) -> CloudOrganization | None:
294    def get_organization(
295        self,
296        *,
297        raise_on_error: bool = True,
298    ) -> CloudOrganization | None:
299        """Get the organization this workspace belongs to.
300
301        Fetching organization info requires ORGANIZATION_READER permissions on the organization,
302        which may not be available with workspace-scoped credentials.
303
304        Args:
305            raise_on_error: If True (default), raises AirbyteError on permission or API errors.
306                If False, returns None instead of raising.
307
308        Returns:
309            CloudOrganization object with organization_id and organization_name,
310            or None if raise_on_error=False and an error occurred.
311
312        Raises:
313            AirbyteError: If raise_on_error=True and the organization info cannot be fetched
314                (e.g., due to insufficient permissions or missing data).
315        """
316        try:
317            info = self._organization_info
318        except (AirbyteError, NotImplementedError):
319            if raise_on_error:
320                raise
321            return None
322
323        organization_id = info.get("organizationId")
324        organization_name = info.get("organizationName")
325
326        # Validate that both organization_id and organization_name are non-null and non-empty
327        if not organization_id or not organization_name:
328            if raise_on_error:
329                raise AirbyteError(
330                    message="Organization info is incomplete.",
331                    context={
332                        "organization_id": organization_id,
333                        "organization_name": organization_name,
334                    },
335                )
336            return None
337
338        organization_credentials = self._credentials.with_organization_id(organization_id)
339        return CloudOrganization(
340            organization_id=organization_id,
341            organization_name=organization_name,
342            client_id=organization_credentials.client_id,
343            client_secret=organization_credentials.client_secret,
344            bearer_token=organization_credentials.bearer_token,
345            public_api_root=organization_credentials.public_api_root,
346            config_api_root=organization_credentials.config_api_root,
347        )

Get the organization this workspace belongs to.

Fetching organization info requires ORGANIZATION_READER permissions on the organization, which may not be available with workspace-scoped credentials.

Arguments:
  • raise_on_error: If True (default), raises AirbyteError on permission or API errors. If False, returns None instead of raising.
Returns:

CloudOrganization object with organization_id and organization_name, or None if raise_on_error=False and an error occurred.

Raises:
  • AirbyteError: If raise_on_error=True and the organization info cannot be fetched (e.g., due to insufficient permissions or missing data).
def connect(self) -> None:
351    def connect(self) -> None:
352        """Check that the workspace is reachable and raise an exception otherwise.
353
354        Note: It is not necessary to call this method before calling other operations. It
355              serves primarily as a simple check to ensure that the workspace is reachable
356              and credentials are correct.
357        """
358        _ = api_util.get_workspace(
359            api_root=self.api_root,
360            workspace_id=self.workspace_id,
361            client_id=self.client_id,
362            client_secret=self.client_secret,
363            bearer_token=self.bearer_token,
364        )
365        print(f"Successfully connected to workspace: {self.workspace_url}")

Check that the workspace is reachable and raise an exception otherwise.

Note: It is not necessary to call this method before calling other operations. It serves primarily as a simple check to ensure that the workspace is reachable and credentials are correct.

def get_connection(self, connection_id: str) -> CloudConnection:
369    def get_connection(
370        self,
371        connection_id: str,
372    ) -> CloudConnection:
373        """Get a connection by ID.
374
375        This method does not fetch data from the API. It returns a `CloudConnection` object,
376        which will be loaded lazily as needed.
377        """
378        return CloudConnection(
379            workspace=self,
380            connection_id=connection_id,
381        )

Get a connection by ID.

This method does not fetch data from the API. It returns a CloudConnection object, which will be loaded lazily as needed.

def get_source(self, source_id: str) -> airbyte.cloud.connectors.CloudSource:
383    def get_source(
384        self,
385        source_id: str,
386    ) -> CloudSource:
387        """Get a source by ID.
388
389        This method does not fetch data from the API. It returns a `CloudSource` object,
390        which will be loaded lazily as needed.
391        """
392        return CloudSource(
393            workspace=self,
394            connector_id=source_id,
395        )

Get a source by ID.

This method does not fetch data from the API. It returns a CloudSource object, which will be loaded lazily as needed.

def get_destination(self, destination_id: str) -> airbyte.cloud.connectors.CloudDestination:
397    def get_destination(
398        self,
399        destination_id: str,
400    ) -> CloudDestination:
401        """Get a destination by ID.
402
403        This method does not fetch data from the API. It returns a `CloudDestination` object,
404        which will be loaded lazily as needed.
405        """
406        return CloudDestination(
407            workspace=self,
408            connector_id=destination_id,
409        )

Get a destination by ID.

This method does not fetch data from the API. It returns a CloudDestination object, which will be loaded lazily as needed.

def check_connector_setup( self, connector_type: Literal['source', 'destination'], connector_id: str) -> airbyte.cloud.models.CheckResult:
411    def check_connector_setup(
412        self,
413        connector_type: Literal["source", "destination"],
414        connector_id: str,
415    ) -> CheckResult:
416        """Run one connection check on a connector that belongs to this workspace.
417
418        Confirms a person has finished a deferred-credential setup in Airbyte Cloud. The
419        connector's workspace is verified first so a check can never be run against a connector
420        outside this workspace.
421        """
422        connector: CloudSource | CloudDestination
423        if connector_type == "source":
424            owner_id = api_util.get_source(
425                source_id=connector_id,
426                api_root=self.api_root,
427                client_id=self.client_id,
428                client_secret=self.client_secret,
429                bearer_token=self.bearer_token,
430            ).workspace_id
431            connector = self.get_source(connector_id)
432        else:
433            owner_id = api_util.get_destination(
434                destination_id=connector_id,
435                api_root=self.api_root,
436                client_id=self.client_id,
437                client_secret=self.client_secret,
438                bearer_token=self.bearer_token,
439            ).workspace_id
440            connector = self.get_destination(connector_id)
441        if owner_id != self.workspace_id:
442            raise exc.AirbyteMissingResourceError(
443                resource_type=connector_type,
444                resource_name_or_id=connector_id,
445                context={"workspace_id": self.workspace_id},
446            )
447        try:
448            return connector.check(raise_on_error=False)
449        except AirbyteError as ex:
450            status_code = (ex.context or {}).get("status_code")
451            if status_code == HTTPStatus.UNPROCESSABLE_ENTITY:
452                return CheckResult(success=False)
453            raise AirbyteError(
454                message="Cloud could not check the connector setup.",
455                context={"status_code": status_code},
456            ) from None

Run one connection check on a connector that belongs to this workspace.

Confirms a person has finished a deferred-credential setup in Airbyte Cloud. The connector's workspace is verified first so a check can never be run against a connector outside this workspace.

def deploy_source( self, name: str, source: airbyte.Source | dict[str, typing.Any], *, unique: bool = True, random_name_suffix: bool = False, definition_id: str | None = None, defer_credentials: bool = False) -> airbyte.cloud.connectors.CloudSource:
460    def deploy_source(
461        self,
462        name: str,
463        source: Source | dict[str, Any],
464        *,
465        unique: bool = True,
466        random_name_suffix: bool = False,
467        definition_id: str | None = None,
468        defer_credentials: bool = False,
469    ) -> CloudSource:
470        """Deploy a source to the workspace.
471
472        Returns the newly deployed source.
473
474        Args:
475            name: The name to use when deploying.
476            source: The source object to deploy, or (with `defer_credentials=True`) a
477                dictionary of non-secret configuration values.
478            unique: Whether to require a unique name. If `True`, duplicate names
479                are not allowed. Defaults to `True`.
480            random_name_suffix: Whether to append a random suffix to the name.
481            definition_id: The source definition ID. Required with `defer_credentials=True`.
482            defer_credentials: Save a draft with partial configuration. A person completes
483                credentials and other missing settings at the returned source's `connector_url`.
484                A successful connection check promotes the draft. Raises
485                `AirbyteDeferredSetupError` if Cloud does not acknowledge draft mode.
486        """
487        if defer_credentials:
488            return CloudSource(
489                workspace=self,
490                connector_id=self._deploy_deferred(
491                    connector_type="source",
492                    name=name,
493                    config=source,
494                    definition_id=definition_id,
495                    unique=unique,
496                    random_name_suffix=random_name_suffix,
497                ),
498            )
499        if isinstance(source, dict):
500            raise exc.PyAirbyteInputError(
501                message="`source` must be a `Source` object unless `defer_credentials=True`.",
502            )
503
504        source_config_dict = source._hydrated_config.copy()  # noqa: SLF001 (non-public API)
505        source_config_dict["sourceType"] = source.name.replace("source-", "")
506
507        if random_name_suffix:
508            name += f" (ID: {text_util.generate_random_suffix()})"
509
510        if unique:
511            existing = self.list_sources(name=name)
512            if existing:
513                raise exc.AirbyteDuplicateResourcesError(
514                    resource_type="source",
515                    resource_name=name,
516                )
517
518        deployed_source = api_util.create_source(
519            name=name,
520            api_root=self.api_root,
521            workspace_id=self.workspace_id,
522            config=source_config_dict,
523            definition_id=definition_id,
524            client_id=self.client_id,
525            client_secret=self.client_secret,
526            bearer_token=self.bearer_token,
527        )
528        return CloudSource(
529            workspace=self,
530            connector_id=deployed_source.source_id,
531        )

Deploy a source to the workspace.

Returns the newly deployed source.

Arguments:
  • name: The name to use when deploying.
  • source: The source object to deploy, or (with defer_credentials=True) a dictionary of non-secret configuration values.
  • unique: Whether to require a unique name. If True, duplicate names are not allowed. Defaults to True.
  • random_name_suffix: Whether to append a random suffix to the name.
  • definition_id: The source definition ID. Required with defer_credentials=True.
  • defer_credentials: Save a draft with partial configuration. A person completes credentials and other missing settings at the returned source's connector_url. A successful connection check promotes the draft. Raises AirbyteDeferredSetupError if Cloud does not acknowledge draft mode.
def deploy_destination( self, name: str, destination: airbyte.Destination | dict[str, typing.Any], *, unique: bool = True, random_name_suffix: bool = False, definition_id: str | None = None, defer_credentials: bool = False) -> airbyte.cloud.connectors.CloudDestination:
533    def deploy_destination(
534        self,
535        name: str,
536        destination: Destination | dict[str, Any],
537        *,
538        unique: bool = True,
539        random_name_suffix: bool = False,
540        definition_id: str | None = None,
541        defer_credentials: bool = False,
542    ) -> CloudDestination:
543        """Deploy a destination to the workspace.
544
545        Returns the newly deployed destination ID.
546
547        Args:
548            name: The name to use when deploying.
549            destination: The destination to deploy. Can be a local Airbyte `Destination` object or a
550                dictionary of configuration values.
551            unique: Whether to require a unique name. If `True`, duplicate names
552                are not allowed. Defaults to `True`.
553            random_name_suffix: Whether to append a random suffix to the name.
554            definition_id: The destination definition ID. Required with `defer_credentials=True`;
555                otherwise the type is inferred from `destinationType`.
556            defer_credentials: Create the destination without its credentials. See
557                `deploy_source`.
558        """
559        if defer_credentials:
560            return CloudDestination(
561                workspace=self,
562                connector_id=self._deploy_deferred(
563                    connector_type="destination",
564                    name=name,
565                    config=destination,
566                    definition_id=definition_id,
567                    unique=unique,
568                    random_name_suffix=random_name_suffix,
569                ),
570            )
571
572        if isinstance(destination, Destination):
573            destination_conf_dict = destination._hydrated_config.copy()  # noqa: SLF001 (non-public API)
574            destination_conf_dict["destinationType"] = destination.name.replace("destination-", "")
575            # raise ValueError(destination_conf_dict)
576        else:
577            destination_conf_dict = destination.copy()
578            if "destinationType" not in destination_conf_dict:
579                raise exc.PyAirbyteInputError(
580                    message="Missing `destinationType` in configuration dictionary.",
581                )
582
583        if random_name_suffix:
584            name += f" (ID: {text_util.generate_random_suffix()})"
585
586        if unique:
587            existing = self.list_destinations(name=name)
588            if existing:
589                raise exc.AirbyteDuplicateResourcesError(
590                    resource_type="destination",
591                    resource_name=name,
592                )
593
594        deployed_destination = api_util.create_destination(
595            name=name,
596            api_root=self.api_root,
597            workspace_id=self.workspace_id,
598            config=destination_conf_dict,  # Wants a dataclass but accepts dict
599            client_id=self.client_id,
600            client_secret=self.client_secret,
601            bearer_token=self.bearer_token,
602        )
603        return CloudDestination(
604            workspace=self,
605            connector_id=deployed_destination.destination_id,
606        )

Deploy a destination to the workspace.

Returns the newly deployed destination ID.

Arguments:
  • name: The name to use when deploying.
  • destination: The destination to deploy. Can be a local Airbyte Destination object or a dictionary of configuration values.
  • unique: Whether to require a unique name. If True, duplicate names are not allowed. Defaults to True.
  • random_name_suffix: Whether to append a random suffix to the name.
  • definition_id: The destination definition ID. Required with defer_credentials=True; otherwise the type is inferred from destinationType.
  • defer_credentials: Create the destination without its credentials. See deploy_source.
def permanently_delete_source( self, source: str | airbyte.cloud.connectors.CloudSource, *, safe_mode: bool = True) -> None:
651    def permanently_delete_source(
652        self,
653        source: str | CloudSource,
654        *,
655        safe_mode: bool = True,
656    ) -> None:
657        """Delete a source from the workspace.
658
659        You can pass either the source ID `str` or a deployed `Source` object.
660
661        Args:
662            source: The source ID or CloudSource object to delete
663            safe_mode: If True, requires the source name to contain "delete-me" or "deleteme"
664                (case insensitive) to prevent accidental deletion. Defaults to True.
665        """
666        if not isinstance(source, (str, CloudSource)):
667            raise exc.PyAirbyteInputError(
668                message="Invalid source type.",
669                input_value=type(source).__name__,
670            )
671
672        api_util.delete_source(
673            source_id=source.connector_id if isinstance(source, CloudSource) else source,
674            source_name=source.name if isinstance(source, CloudSource) else None,
675            api_root=self.api_root,
676            client_id=self.client_id,
677            client_secret=self.client_secret,
678            bearer_token=self.bearer_token,
679            safe_mode=safe_mode,
680        )

Delete a source from the workspace.

You can pass either the source ID str or a deployed Source object.

Arguments:
  • source: The source ID or CloudSource object to delete
  • safe_mode: If True, requires the source name to contain "delete-me" or "deleteme" (case insensitive) to prevent accidental deletion. Defaults to True.
def permanently_delete_destination( self, destination: str | airbyte.cloud.connectors.CloudDestination, *, safe_mode: bool = True) -> None:
684    def permanently_delete_destination(
685        self,
686        destination: str | CloudDestination,
687        *,
688        safe_mode: bool = True,
689    ) -> None:
690        """Delete a deployed destination from the workspace.
691
692        You can pass either the `Cache` class or the deployed destination ID as a `str`.
693
694        Args:
695            destination: The destination ID or CloudDestination object to delete
696            safe_mode: If True, requires the destination name to contain "delete-me" or "deleteme"
697                (case insensitive) to prevent accidental deletion. Defaults to True.
698        """
699        if not isinstance(destination, (str, CloudDestination)):
700            raise exc.PyAirbyteInputError(
701                message="Invalid destination type.",
702                input_value=type(destination).__name__,
703            )
704
705        api_util.delete_destination(
706            destination_id=(
707                destination if isinstance(destination, str) else destination.destination_id
708            ),
709            destination_name=(
710                destination.name if isinstance(destination, CloudDestination) else None
711            ),
712            api_root=self.api_root,
713            client_id=self.client_id,
714            client_secret=self.client_secret,
715            bearer_token=self.bearer_token,
716            safe_mode=safe_mode,
717        )

Delete a deployed destination from the workspace.

You can pass either the Cache class or the deployed destination ID as a str.

Arguments:
  • destination: The destination ID or CloudDestination object to delete
  • safe_mode: If True, requires the destination name to contain "delete-me" or "deleteme" (case insensitive) to prevent accidental deletion. Defaults to True.
def deploy_connection( self, connection_name: str, *, source: airbyte.cloud.connectors.CloudSource | str, selected_streams: list[str], destination: airbyte.cloud.connectors.CloudDestination | str, table_prefix: str | None = None) -> CloudConnection:
721    def deploy_connection(
722        self,
723        connection_name: str,
724        *,
725        source: CloudSource | str,
726        selected_streams: list[str],
727        destination: CloudDestination | str,
728        table_prefix: str | None = None,
729    ) -> CloudConnection:
730        """Create a new connection between an already deployed source and destination.
731
732        Returns the newly deployed connection object.
733
734        Args:
735            connection_name: The name of the connection.
736            source: The deployed source. You can pass a source ID or a CloudSource object.
737            destination: The deployed destination. You can pass a destination ID or a
738                CloudDestination object.
739            table_prefix: Optional. The table prefix to use when syncing to the destination.
740            selected_streams: The selected stream names to sync within the connection.
741        """
742        if not selected_streams:
743            raise exc.PyAirbyteInputError(
744                guidance="You must provide `selected_streams` when creating a connection."
745            )
746
747        source_id: str = source if isinstance(source, str) else source.connector_id
748        destination_id: str = (
749            destination if isinstance(destination, str) else destination.connector_id
750        )
751
752        deployed_connection = api_util.create_connection(
753            name=connection_name,
754            source_id=source_id,
755            destination_id=destination_id,
756            api_root=self.api_root,
757            workspace_id=self.workspace_id,
758            selected_stream_names=selected_streams,
759            prefix=table_prefix or "",
760            client_id=self.client_id,
761            client_secret=self.client_secret,
762            bearer_token=self.bearer_token,
763        )
764
765        return CloudConnection(
766            workspace=self,
767            connection_id=deployed_connection.connection_id,
768            source=deployed_connection.source_id,
769            destination=deployed_connection.destination_id,
770        )

Create a new connection between an already deployed source and destination.

Returns the newly deployed connection object.

Arguments:
  • connection_name: The name of the connection.
  • source: The deployed source. You can pass a source ID or a CloudSource object.
  • destination: The deployed destination. You can pass a destination ID or a CloudDestination object.
  • table_prefix: Optional. The table prefix to use when syncing to the destination.
  • selected_streams: The selected stream names to sync within the connection.
def permanently_delete_connection( self, connection: str | CloudConnection, *, cascade_delete_source: bool = False, cascade_delete_destination: bool = False, safe_mode: bool = True) -> None:
772    def permanently_delete_connection(
773        self,
774        connection: str | CloudConnection,
775        *,
776        cascade_delete_source: bool = False,
777        cascade_delete_destination: bool = False,
778        safe_mode: bool = True,
779    ) -> None:
780        """Delete a deployed connection from the workspace.
781
782        Args:
783            connection: The connection ID or CloudConnection object to delete
784            cascade_delete_source: If True, also delete the source after deleting the connection
785            cascade_delete_destination: If True, also delete the destination after deleting
786                the connection
787            safe_mode: If True, requires the connection name to contain "delete-me" or "deleteme"
788                (case insensitive) to prevent accidental deletion. Defaults to True. Also applies
789                to cascade deletes.
790        """
791        if connection is None:
792            raise ValueError("No connection ID provided.")
793
794        if isinstance(connection, str):
795            connection = CloudConnection(
796                workspace=self,
797                connection_id=connection,
798            )
799
800        api_util.delete_connection(
801            connection_id=connection.connection_id,
802            connection_name=connection.name,
803            api_root=self.api_root,
804            workspace_id=self.workspace_id,
805            client_id=self.client_id,
806            client_secret=self.client_secret,
807            bearer_token=self.bearer_token,
808            safe_mode=safe_mode,
809        )
810
811        if cascade_delete_source:
812            self.permanently_delete_source(
813                source=connection.source_id,
814                safe_mode=safe_mode,
815            )
816        if cascade_delete_destination:
817            self.permanently_delete_destination(
818                destination=connection.destination_id,
819                safe_mode=safe_mode,
820            )

Delete a deployed connection from the workspace.

Arguments:
  • connection: The connection ID or CloudConnection object to delete
  • cascade_delete_source: If True, also delete the source after deleting the connection
  • cascade_delete_destination: If True, also delete the destination after deleting the connection
  • safe_mode: If True, requires the connection name to contain "delete-me" or "deleteme" (case insensitive) to prevent accidental deletion. Defaults to True. Also applies to cascade deletes.
def list_workspaces( self, name: str | None = None, *, name_filter: Callable | None = None, limit: int | None = None) -> list[CloudWorkspaceInfo]:
824    def list_workspaces(
825        self,
826        name: str | None = None,
827        *,
828        name_filter: Callable | None = None,
829        limit: int | None = None,
830    ) -> list[CloudWorkspaceInfo]:
831        """List workspaces available to the current credentials, with an optional limit."""
832        return [
833            CloudWorkspaceInfo.from_api_response(workspace)
834            for workspace in api_util.list_workspaces(
835                workspace_id="",
836                api_root=self.api_root,
837                name=name,
838                name_filter=name_filter,
839                client_id=self.client_id,
840                client_secret=self.client_secret,
841                bearer_token=self.bearer_token,
842                limit=limit,
843            )
844        ]

List workspaces available to the current credentials, with an optional limit.

def rename(self, name: str) -> CloudWorkspace:
846    def rename(
847        self,
848        name: str,
849    ) -> CloudWorkspace:
850        """Rename this workspace."""
851        api_util.rename_workspace(
852            workspace_id=self.workspace_id,
853            name=name,
854            api_root=self.api_root,
855            client_id=self.client_id,
856            client_secret=self.client_secret,
857            bearer_token=self.bearer_token,
858        )
859        return self

Rename this workspace.

def permanently_delete( self, *, workspace_name: str | None = None, safe_mode: bool = True) -> None:
861    def permanently_delete(
862        self,
863        *,
864        workspace_name: str | None = None,
865        safe_mode: bool = True,
866    ) -> None:
867        """Permanently delete this workspace if it has no connections.
868
869        When `safe_mode` is enabled, the workspace name must contain `delete-me`
870        or `deleteme`. This also checks for existing connections before deleting
871        and raises `AirbyteWorkspaceNotEmptyError` if the workspace is not empty.
872        """
873        api_util.permanently_delete_workspace(
874            workspace_id=self.workspace_id,
875            workspace_name=workspace_name,
876            api_root=self.api_root,
877            client_id=self.client_id,
878            client_secret=self.client_secret,
879            bearer_token=self.bearer_token,
880            safe_mode=safe_mode,
881        )

Permanently delete this 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_connections( self, name: str | None = None, *, name_filter: Callable | None = None, limit: int | None = None) -> list[CloudConnection]:
883    def list_connections(
884        self,
885        name: str | None = None,
886        *,
887        name_filter: Callable | None = None,
888        limit: int | None = None,
889    ) -> list[CloudConnection]:
890        """List connections by name in the workspace, with an optional limit."""
891        connections = api_util.list_connections(
892            api_root=self.api_root,
893            workspace_id=self.workspace_id,
894            name=name,
895            name_filter=name_filter,
896            limit=limit,
897            client_id=self.client_id,
898            client_secret=self.client_secret,
899            bearer_token=self.bearer_token,
900        )
901        return [
902            CloudConnection._from_connection_response(  # noqa: SLF001 (non-public API)
903                workspace=self,
904                connection_response=connection,
905            )
906            for connection in connections
907        ]

List connections by name in the workspace, with an optional limit.

def list_sources( self, name: str | None = None, *, name_filter: Callable | None = None, limit: int | None = None) -> list[airbyte.cloud.connectors.CloudSource]:
909    def list_sources(
910        self,
911        name: str | None = None,
912        *,
913        name_filter: Callable | None = None,
914        limit: int | None = None,
915    ) -> list[CloudSource]:
916        """List all sources in the workspace, with an optional limit."""
917        sources = api_util.list_sources(
918            api_root=self.api_root,
919            workspace_id=self.workspace_id,
920            name=name,
921            name_filter=name_filter,
922            limit=limit,
923            client_id=self.client_id,
924            client_secret=self.client_secret,
925            bearer_token=self.bearer_token,
926        )
927        return [
928            CloudSource._from_source_response(  # noqa: SLF001 (non-public API)
929                workspace=self,
930                source_response=source,
931            )
932            for source in sources
933        ]

List all sources in the workspace, with an optional limit.

def list_destinations( self, name: str | None = None, *, name_filter: Callable | None = None, limit: int | None = None) -> list[airbyte.cloud.connectors.CloudDestination]:
935    def list_destinations(
936        self,
937        name: str | None = None,
938        *,
939        name_filter: Callable | None = None,
940        limit: int | None = None,
941    ) -> list[CloudDestination]:
942        """List all destinations in the workspace, with an optional limit."""
943        destinations = api_util.list_destinations(
944            api_root=self.api_root,
945            workspace_id=self.workspace_id,
946            name=name,
947            name_filter=name_filter,
948            limit=limit,
949            client_id=self.client_id,
950            client_secret=self.client_secret,
951            bearer_token=self.bearer_token,
952        )
953        return [
954            CloudDestination._from_destination_response(  # noqa: SLF001 (non-public API)
955                workspace=self,
956                destination_response=destination,
957            )
958            for destination in destinations
959        ]

List all destinations in the workspace, with an optional limit.

def publish_custom_source_definition( self, name: str, *, manifest_yaml: dict[str, typing.Any] | pathlib.Path | str | None = None, docker_image: str | None = None, docker_tag: str | None = None, unique: bool = True, pre_validate: bool = True, testing_values: dict[str, typing.Any] | None = None) -> airbyte.cloud.connectors.CustomCloudSourceDefinition:
 961    def publish_custom_source_definition(
 962        self,
 963        name: str,
 964        *,
 965        manifest_yaml: dict[str, Any] | Path | str | None = None,
 966        docker_image: str | None = None,
 967        docker_tag: str | None = None,
 968        unique: bool = True,
 969        pre_validate: bool = True,
 970        testing_values: dict[str, Any] | None = None,
 971    ) -> CustomCloudSourceDefinition:
 972        """Publish a custom source connector definition.
 973
 974        You must specify EITHER manifest_yaml (for YAML connectors) OR both docker_image
 975        and docker_tag (for Docker connectors), but not both.
 976
 977        Args:
 978            name: Display name for the connector definition
 979            manifest_yaml: Low-code CDK manifest (dict, Path to YAML file, or YAML string)
 980            docker_image: Docker repository (e.g., 'airbyte/source-custom')
 981            docker_tag: Docker image tag (e.g., '1.0.0')
 982            unique: Whether to enforce name uniqueness
 983            pre_validate: Whether to validate manifest client-side (YAML only)
 984            testing_values: Optional configuration values to use for testing in the
 985                Connector Builder UI. If provided, these values are stored as the complete
 986                testing values object for the connector builder project (replaces any existing
 987                values), allowing immediate test read operations.
 988
 989        Returns:
 990            CustomCloudSourceDefinition object representing the created definition
 991
 992        Raises:
 993            PyAirbyteInputError: If both or neither of manifest_yaml and docker_image provided
 994            AirbyteDuplicateResourcesError: If unique=True and name already exists
 995        """
 996        is_yaml = manifest_yaml is not None
 997        is_docker = docker_image is not None
 998
 999        if is_yaml == is_docker:
1000            raise exc.PyAirbyteInputError(
1001                message=(
1002                    "Must specify EITHER manifest_yaml (for YAML connectors) OR "
1003                    "docker_image + docker_tag (for Docker connectors), but not both"
1004                ),
1005                context={
1006                    "manifest_yaml_provided": is_yaml,
1007                    "docker_image_provided": is_docker,
1008                },
1009            )
1010
1011        if is_docker and docker_tag is None:
1012            raise exc.PyAirbyteInputError(
1013                message="docker_tag is required when docker_image is specified",
1014                context={"docker_image": docker_image},
1015            )
1016
1017        if unique:
1018            existing = self.list_custom_source_definitions(
1019                definition_type="yaml" if is_yaml else "docker",
1020            )
1021            if any(d.name == name for d in existing):
1022                raise exc.AirbyteDuplicateResourcesError(
1023                    resource_type="custom_source_definition",
1024                    resource_name=name,
1025                )
1026
1027        if is_yaml:
1028            manifest_dict: dict[str, Any]
1029            if isinstance(manifest_yaml, Path):
1030                manifest_dict = yaml.safe_load(manifest_yaml.read_text())
1031            elif isinstance(manifest_yaml, str):
1032                manifest_dict = yaml.safe_load(manifest_yaml)
1033            elif manifest_yaml is not None:
1034                manifest_dict = manifest_yaml
1035            else:
1036                raise exc.PyAirbyteInputError(
1037                    message="manifest_yaml is required for YAML connectors",
1038                    context={"name": name},
1039                )
1040
1041            if pre_validate:
1042                api_util.validate_yaml_manifest(manifest_dict, raise_on_error=True)
1043
1044            result = api_util.create_custom_yaml_source_definition(
1045                name=name,
1046                workspace_id=self.workspace_id,
1047                manifest=manifest_dict,
1048                api_root=self.api_root,
1049                client_id=self.client_id,
1050                client_secret=self.client_secret,
1051                bearer_token=self.bearer_token,
1052            )
1053            custom_definition = CustomCloudSourceDefinition._from_yaml_response(  # noqa: SLF001
1054                self, result
1055            )
1056
1057            # Set testing values if provided
1058            if testing_values is not None:
1059                custom_definition.set_testing_values(testing_values)
1060
1061            return custom_definition
1062
1063        raise NotImplementedError(
1064            "Docker custom source definitions are not yet supported. "
1065            "Only YAML manifest-based custom sources are currently available."
1066        )

Publish a custom source connector definition.

You must specify EITHER manifest_yaml (for YAML connectors) OR both docker_image and docker_tag (for Docker connectors), but not both.

Arguments:
  • name: Display name for the connector definition
  • manifest_yaml: Low-code CDK manifest (dict, Path to YAML file, or YAML string)
  • docker_image: Docker repository (e.g., 'airbyte/source-custom')
  • docker_tag: Docker image tag (e.g., '1.0.0')
  • unique: Whether to enforce name uniqueness
  • pre_validate: Whether to validate manifest client-side (YAML only)
  • testing_values: Optional configuration values to use for testing in the Connector Builder UI. If provided, these values are stored as the complete testing values object for the connector builder project (replaces any existing values), allowing immediate test read operations.
Returns:

CustomCloudSourceDefinition object representing the created definition

Raises:
  • PyAirbyteInputError: If both or neither of manifest_yaml and docker_image provided
  • AirbyteDuplicateResourcesError: If unique=True and name already exists
def list_custom_source_definitions( self, *, definition_type: Literal['yaml', 'docker']) -> list[airbyte.cloud.connectors.CustomCloudSourceDefinition]:
1068    def list_custom_source_definitions(
1069        self,
1070        *,
1071        definition_type: Literal["yaml", "docker"],
1072    ) -> list[CustomCloudSourceDefinition]:
1073        """List custom source connector definitions.
1074
1075        Args:
1076            definition_type: Connector type to list ("yaml" or "docker"). Required.
1077
1078        Returns:
1079            List of CustomCloudSourceDefinition objects matching the specified type
1080        """
1081        if definition_type == "yaml":
1082            yaml_definitions = api_util.list_custom_yaml_source_definitions(
1083                workspace_id=self.workspace_id,
1084                api_root=self.api_root,
1085                client_id=self.client_id,
1086                client_secret=self.client_secret,
1087                bearer_token=self.bearer_token,
1088            )
1089            return [
1090                CustomCloudSourceDefinition._from_yaml_response(self, d)  # noqa: SLF001
1091                for d in yaml_definitions
1092            ]
1093
1094        raise NotImplementedError(
1095            "Docker custom source definitions are not yet supported. "
1096            "Only YAML manifest-based custom sources are currently available."
1097        )

List custom source connector definitions.

Arguments:
  • definition_type: Connector type to list ("yaml" or "docker"). Required.
Returns:

List of CustomCloudSourceDefinition objects matching the specified type

def get_custom_source_definition( self, definition_id: str, *, definition_type: Literal['yaml', 'docker']) -> airbyte.cloud.connectors.CustomCloudSourceDefinition:
1099    def get_custom_source_definition(
1100        self,
1101        definition_id: str,
1102        *,
1103        definition_type: Literal["yaml", "docker"],
1104    ) -> CustomCloudSourceDefinition:
1105        """Get a specific custom source definition by ID.
1106
1107        Args:
1108            definition_id: The definition ID
1109            definition_type: Connector type ("yaml" or "docker"). Required.
1110
1111        Returns:
1112            CustomCloudSourceDefinition object
1113        """
1114        if definition_type == "yaml":
1115            result = api_util.get_custom_yaml_source_definition(
1116                workspace_id=self.workspace_id,
1117                definition_id=definition_id,
1118                api_root=self.api_root,
1119                client_id=self.client_id,
1120                client_secret=self.client_secret,
1121                bearer_token=self.bearer_token,
1122            )
1123            return CustomCloudSourceDefinition._from_yaml_response(self, result)  # noqa: SLF001
1124
1125        raise NotImplementedError(
1126            "Docker custom source definitions are not yet supported. "
1127            "Only YAML manifest-based custom sources are currently available."
1128        )

Get a specific custom source definition by ID.

Arguments:
  • definition_id: The definition ID
  • definition_type: Connector type ("yaml" or "docker"). Required.
Returns:

CustomCloudSourceDefinition object

class CloudConnection:
  77class CloudConnection:  # noqa: PLR0904  # Too many public methods
  78    """A connection is an extract-load (EL) pairing of a source and destination in Airbyte Cloud.
  79
  80    You can use a connection object to run sync jobs, retrieve logs, and manage the connection.
  81    """
  82
  83    def __init__(
  84        self,
  85        workspace: CloudWorkspace,
  86        connection_id: str,
  87        source: str | None = None,
  88        destination: str | None = None,
  89    ) -> None:
  90        """It is not recommended to create a `CloudConnection` object directly.
  91
  92        Instead, use `CloudWorkspace.get_connection()` to create a connection object.
  93        """
  94        self.connection_id = connection_id
  95        """The ID of the connection."""
  96
  97        self.workspace = workspace
  98        """The workspace that the connection belongs to."""
  99
 100        self._source_id = source
 101        """The ID of the source."""
 102
 103        self._destination_id = destination
 104        """The ID of the destination."""
 105
 106        self._connection_info: CloudConnectionInfo | None = None
 107        """The connection info object. (Cached.)"""
 108
 109        self._cloud_source_object: CloudSource | None = None
 110        """The source object. (Cached.)"""
 111
 112        self._cloud_destination_object: CloudDestination | None = None
 113        """The destination object. (Cached.)"""
 114
 115    def _fetch_connection_info(
 116        self,
 117        *,
 118        force_refresh: bool = False,
 119        verify: bool = True,
 120    ) -> CloudConnectionInfo:
 121        """Fetch and cache connection info from the API.
 122
 123        By default, this method will only fetch from the API if connection info is not
 124        already cached. It also verifies that the connection belongs to the expected
 125        workspace unless verification is explicitly disabled.
 126
 127        Args:
 128            force_refresh: If True, always fetch from the API even if cached.
 129                If False (default), only fetch if not already cached.
 130            verify: If True (default), verify that the connection is valid (e.g., that
 131                the workspace_id matches this object's workspace). Raises an error if
 132                validation fails.
 133
 134        Returns:
 135            Information about the connection from the API.
 136
 137        Raises:
 138            AirbyteWorkspaceMismatchError: If verify is True and the connection's
 139                workspace_id doesn't match the expected workspace.
 140            AirbyteMissingResourceError: If the connection doesn't exist.
 141        """
 142        if not force_refresh and self._connection_info is not None:
 143            # Use cached info, but still verify if requested
 144            if verify:
 145                self._verify_workspace_match(self._connection_info)
 146            return self._connection_info
 147
 148        # Fetch from API
 149        connection_info = api_util.get_connection(
 150            workspace_id=self.workspace.workspace_id,
 151            connection_id=self.connection_id,
 152            api_root=self.workspace.api_root,
 153            client_id=self.workspace.client_id,
 154            client_secret=self.workspace.client_secret,
 155            bearer_token=self.workspace.bearer_token,
 156        )
 157        result = CloudConnectionInfo.from_api_response(connection_info)
 158
 159        self._connection_info = result
 160
 161        # Verify if requested
 162        if verify:
 163            self._verify_workspace_match(result)
 164
 165        return result
 166
 167    def _verify_workspace_match(self, connection_info: CloudConnectionInfo) -> None:
 168        """Verify that the connection belongs to the expected workspace.
 169
 170        Raises:
 171            AirbyteWorkspaceMismatchError: If the workspace IDs don't match.
 172        """
 173        if connection_info.workspace_id != self.workspace.workspace_id:
 174            raise AirbyteWorkspaceMismatchError(
 175                resource_type="connection",
 176                resource_id=self.connection_id,
 177                workspace=self.workspace,
 178                expected_workspace_id=self.workspace.workspace_id,
 179                actual_workspace_id=connection_info.workspace_id,
 180                message=(
 181                    f"Connection '{self.connection_id}' belongs to workspace "
 182                    f"'{connection_info.workspace_id}', not '{self.workspace.workspace_id}'."
 183                ),
 184            )
 185
 186    def check_is_valid(self) -> bool:
 187        """Check if this connection exists and belongs to the expected workspace.
 188
 189        This method fetches connection info from the API (if not already cached) and
 190        verifies that the connection's workspace_id matches the workspace associated
 191        with this CloudConnection object.
 192
 193        Returns:
 194            True if the connection exists and belongs to the expected workspace.
 195
 196        Raises:
 197            AirbyteWorkspaceMismatchError: If the connection belongs to a different workspace.
 198            AirbyteMissingResourceError: If the connection doesn't exist.
 199        """
 200        self._fetch_connection_info(force_refresh=False, verify=True)
 201        return True
 202
 203    @classmethod
 204    def _from_connection_response(
 205        cls,
 206        workspace: CloudWorkspace,
 207        connection_response: _ConnectionResponseLike,
 208    ) -> CloudConnection:
 209        """Create a CloudConnection from an API connection response."""
 210        connection_info = CloudConnectionInfo.from_api_response(connection_response)
 211        result = cls(
 212            workspace=workspace,
 213            connection_id=connection_info.connection_id,
 214            source=connection_info.source_id,
 215            destination=connection_info.destination_id,
 216        )
 217        result._connection_info = connection_info  # noqa: SLF001 # Accessing Non-Public API
 218        return result
 219
 220    # Properties
 221
 222    @property
 223    def name(self) -> str | None:
 224        """Get the display name of the connection, if available.
 225
 226        E.g. "My Postgres to Snowflake", not the connection ID.
 227        """
 228        if not self._connection_info:
 229            self._connection_info = self._fetch_connection_info()
 230
 231        return self._connection_info.name
 232
 233    @property
 234    def source_id(self) -> str:
 235        """The ID of the source."""
 236        if not self._source_id:
 237            if not self._connection_info:
 238                self._connection_info = self._fetch_connection_info()
 239
 240            self._source_id = self._connection_info.source_id
 241
 242        return self._source_id
 243
 244    @property
 245    def source(self) -> CloudSource:
 246        """Get the source object."""
 247        if self._cloud_source_object:
 248            return self._cloud_source_object
 249
 250        self._cloud_source_object = CloudSource(
 251            workspace=self.workspace,
 252            connector_id=self.source_id,
 253        )
 254        return self._cloud_source_object
 255
 256    @property
 257    def destination_id(self) -> str:
 258        """The ID of the destination."""
 259        if not self._destination_id:
 260            if not self._connection_info:
 261                self._connection_info = self._fetch_connection_info()
 262
 263            self._destination_id = self._connection_info.destination_id
 264
 265        return self._destination_id
 266
 267    @property
 268    def destination(self) -> CloudDestination:
 269        """Get the destination object."""
 270        if self._cloud_destination_object:
 271            return self._cloud_destination_object
 272
 273        self._cloud_destination_object = CloudDestination(
 274            workspace=self.workspace,
 275            connector_id=self.destination_id,
 276        )
 277        return self._cloud_destination_object
 278
 279    @property
 280    def stream_names(self) -> list[str]:
 281        """The stream names."""
 282        if not self._connection_info:
 283            self._connection_info = self._fetch_connection_info()
 284
 285        return [stream.name for stream in self._connection_info.configurations.streams or []]
 286
 287    @property
 288    def table_prefix(self) -> str:
 289        """The table prefix."""
 290        if not self._connection_info:
 291            self._connection_info = self._fetch_connection_info()
 292
 293        return self._connection_info.prefix or ""
 294
 295    @property
 296    def namespace_definition(self) -> str | None:
 297        """How destination namespaces are chosen: `source`, `destination`, or `custom_format`."""
 298        if not self._connection_info:
 299            self._connection_info = self._fetch_connection_info()
 300
 301        return self._connection_info.namespace_definition
 302
 303    @property
 304    def namespace_format(self) -> str | None:
 305        """The namespace format template, when `namespace_definition` is `custom_format`."""
 306        if not self._connection_info:
 307            self._connection_info = self._fetch_connection_info()
 308
 309        return self._connection_info.namespace_format
 310
 311    @property
 312    def connection_url(self) -> str | None:
 313        """The web URL to the connection."""
 314        return f"{self.workspace.workspace_url}/connections/{self.connection_id}"
 315
 316    @property
 317    def job_history_url(self) -> str | None:
 318        """The URL to the job history for the connection."""
 319        return f"{self.connection_url}/timeline"
 320
 321    # Run Sync
 322
 323    def run_sync(
 324        self,
 325        *,
 326        wait: bool = True,
 327        wait_timeout: int = 300,
 328    ) -> SyncResult:
 329        """Run a sync."""
 330        try:
 331            connection_response = api_util.run_connection(
 332                connection_id=self.connection_id,
 333                api_root=self.workspace.api_root,
 334                workspace_id=self.workspace.workspace_id,
 335                client_id=self.workspace.client_id,
 336                client_secret=self.workspace.client_secret,
 337                bearer_token=self.workspace.bearer_token,
 338            )
 339        except AirbyteConnectionSyncError as ex:
 340            if (
 341                ex.context
 342                and ex.context.get("status_code") == HTTPStatus.CONFLICT
 343                and not self.enabled
 344            ):
 345                raise PyAirbyteInputError(
 346                    message=(
 347                        f"Connection '{self.connection_id}' is disabled (status 'inactive'), "
 348                        "so a sync cannot be started."
 349                    ),
 350                    guidance=(
 351                        "Re-enable the connection first (e.g. "
 352                        "`connection.set_enabled(enabled=True)`, or the "
 353                        "`update_cloud_connection` MCP tool with `enabled=True`), then retry."
 354                    ),
 355                    context={"connection_id": self.connection_id},
 356                ) from ex
 357            raise
 358        sync_result = SyncResult(
 359            workspace=self.workspace,
 360            connection=self,
 361            job_id=connection_response.job_id,
 362        )
 363
 364        if wait:
 365            sync_result.wait_for_completion(
 366                wait_timeout=wait_timeout,
 367                raise_failure=True,
 368                raise_timeout=True,
 369            )
 370
 371        return sync_result
 372
 373    def _get_latest_cancellable_sync_job_id(self) -> int:
 374        """Get the latest cancellable sync job ID."""
 375        sync_results = self.get_previous_sync_logs(
 376            limit=1,
 377            job_type=JobTypeEnum.SYNC,
 378        )
 379        sync_result = sync_results[0] if sync_results else None
 380        if sync_result is None:
 381            raise PyAirbyteInputError(
 382                message="No sync jobs found for this connection.",
 383            )
 384        if sync_result.is_job_complete():
 385            raise PyAirbyteInputError(
 386                message=(
 387                    f"The latest sync job is already finished with status "
 388                    f"'{sync_result.get_job_status().value}'. "
 389                    "Pass an explicit job_id to target a different job."
 390                ),
 391            )
 392        return sync_result.job_id
 393
 394    def _validated_cancellable_job_id(self, job_id: int) -> int:
 395        """Validate an explicit cancellable job ID."""
 396        job_info = api_util.get_job_info(
 397            job_id=job_id,
 398            api_root=self.workspace.api_root,
 399            client_id=self.workspace.client_id,
 400            client_secret=self.workspace.client_secret,
 401            bearer_token=self.workspace.bearer_token,
 402        )
 403        if job_info.connection_id != self.connection_id:
 404            raise PyAirbyteInputError(
 405                message=(
 406                    f"Job {job_id} belongs to connection '{job_info.connection_id}', "
 407                    f"not '{self.connection_id}'."
 408                ),
 409            )
 410        job_status = CloudJobInfo.from_api_response(job_info).status
 411        if job_status in FINAL_STATUSES:
 412            raise PyAirbyteInputError(
 413                message=f"Job {job_id} is already finished with status " f"'{job_status.value}'.",
 414            )
 415        return job_id
 416
 417    def cancel_sync(self, job_id: int | None = None) -> SyncResult:
 418        """Cancel a running sync job.
 419
 420        Defaults to the connection's most recent sync job. Other job types must be
 421        targeted with an explicit `job_id`.
 422        """
 423        target_job_id: int = (
 424            self._get_latest_cancellable_sync_job_id()
 425            if job_id is None
 426            else self._validated_cancellable_job_id(job_id)
 427        )
 428
 429        job_response = api_util.cancel_job(
 430            job_id=target_job_id,
 431            api_root=self.workspace.api_root,
 432            client_id=self.workspace.client_id,
 433            client_secret=self.workspace.client_secret,
 434            bearer_token=self.workspace.bearer_token,
 435        )
 436        return SyncResult(
 437            workspace=self.workspace,
 438            connection=self,
 439            job_id=job_response.job_id,
 440            _latest_job_info=CloudJobInfo.from_api_response(job_response),
 441        )
 442
 443    def __repr__(self) -> str:
 444        """String representation of the connection."""
 445        return (
 446            f"CloudConnection(connection_id={self.connection_id}, source_id={self.source_id}, "
 447            f"destination_id={self.destination_id}, connection_url={self.connection_url})"
 448        )
 449
 450    # Logs
 451
 452    def get_previous_sync_logs(
 453        self,
 454        *,
 455        limit: int = 20,
 456        offset: int | None = None,
 457        from_tail: bool = True,
 458        job_type: str | JobTypeEnum | None = None,
 459    ) -> list[SyncResult]:
 460        """Get previous sync jobs for a connection with pagination support.
 461
 462        Returns SyncResult objects containing job metadata (job_id, status, bytes_synced,
 463        rows_synced, start_time). Full log text can be fetched lazily via
 464        `SyncResult.get_full_log_text()`.
 465
 466        Args:
 467            limit: Maximum number of jobs to return. Defaults to 20.
 468            offset: Number of jobs to skip from the beginning. Defaults to None (0).
 469            from_tail: If True, returns jobs ordered newest-first (createdAt DESC).
 470                If False, returns jobs ordered oldest-first (createdAt ASC).
 471                Defaults to True.
 472            job_type: Filter by job type (e.g., `sync`, `refresh`).
 473                If not specified, defaults to sync and reset jobs only (API default behavior).
 474
 475        Returns:
 476            A list of SyncResult objects representing the sync jobs.
 477        """
 478        order_by = (
 479            api_util.JOB_ORDER_BY_CREATED_AT_DESC
 480            if from_tail
 481            else api_util.JOB_ORDER_BY_CREATED_AT_ASC
 482        )
 483        sync_logs = api_util.get_job_logs(
 484            connection_id=self.connection_id,
 485            api_root=self.workspace.api_root,
 486            workspace_id=self.workspace.workspace_id,
 487            limit=limit,
 488            offset=offset,
 489            order_by=order_by,
 490            job_type=job_type,
 491            client_id=self.workspace.client_id,
 492            client_secret=self.workspace.client_secret,
 493            bearer_token=self.workspace.bearer_token,
 494        )
 495        return [
 496            SyncResult(
 497                workspace=self.workspace,
 498                connection=self,
 499                job_id=sync_log.job_id,
 500                _latest_job_info=CloudJobInfo.from_api_response(sync_log),
 501            )
 502            for sync_log in sync_logs
 503        ]
 504
 505    def get_sync_result(
 506        self,
 507        job_id: int | None = None,
 508    ) -> SyncResult | None:
 509        """Get the sync result for the connection.
 510
 511        If `job_id` is not provided, the most recent sync job will be used.
 512
 513        Returns `None` if job_id is omitted and no previous jobs are found.
 514        """
 515        if job_id is None:
 516            # Get the most recent sync job
 517            results = self.get_previous_sync_logs(
 518                limit=1,
 519            )
 520            if results:
 521                return results[0]
 522
 523            return None
 524
 525        # Get the sync job by ID (lazy loaded)
 526        return SyncResult(
 527            workspace=self.workspace,
 528            connection=self,
 529            job_id=job_id,
 530        )
 531
 532    # Artifacts
 533
 534    @deprecated("Use 'dump_raw_state()' instead.")
 535    def get_state_artifacts(self) -> list[dict[str, Any]] | None:
 536        """Deprecated. Use `dump_raw_state()` instead."""
 537        state_response = api_util.get_connection_state(
 538            connection_id=self.connection_id,
 539            api_root=self.workspace.api_root,
 540            client_id=self.workspace.client_id,
 541            client_secret=self.workspace.client_secret,
 542            bearer_token=self.workspace.bearer_token,
 543            config_api_root=self.workspace.config_api_root,
 544        )
 545        if state_response.get("stateType") == "not_set":
 546            return None
 547        return state_response.get("streamState", [])
 548
 549    @overload
 550    def dump_raw_state(self, *, normalize: Literal[True] = True) -> list[dict[str, Any]]: ...
 551
 552    @overload
 553    def dump_raw_state(self, *, normalize: Literal[False]) -> dict[str, Any]: ...
 554
 555    def dump_raw_state(
 556        self,
 557        *,
 558        normalize: bool = True,
 559    ) -> dict[str, Any] | list[dict[str, Any]]:
 560        """Dump the state for this connection.
 561
 562        By default, returns a list of Airbyte protocol `AirbyteStateMessage` dicts
 563        with snake_case keys, suitable for passing to a connector's `--state` flag.
 564
 565        When `normalize` is `False`, returns the raw Config API dict (camelCase keys,
 566        includes `stateType` and `connectionId`). This raw format can be passed
 567        directly to `import_raw_state()` for backup/restore workflows.
 568
 569        Args:
 570            normalize: If `True` (default), convert to Airbyte protocol format.
 571                If `False`, return the raw Config API response.
 572
 573        Returns:
 574            Normalized: list of protocol-format state message dicts (empty list if
 575            no state). Raw: the full Config API state dict.
 576        """
 577        raw = api_util.get_connection_state(
 578            connection_id=self.connection_id,
 579            api_root=self.workspace.api_root,
 580            client_id=self.workspace.client_id,
 581            client_secret=self.workspace.client_secret,
 582            bearer_token=self.workspace.bearer_token,
 583            config_api_root=self.workspace.config_api_root,
 584        )
 585        if normalize:
 586            return _normalize_state_to_protocol(raw)
 587        return raw
 588
 589    def import_raw_state(
 590        self,
 591        connection_state: dict[str, Any] | list[dict[str, Any]],
 592    ) -> dict[str, Any]:
 593        """Import (restore) the full state for this connection.
 594
 595        > ⚠️ **WARNING:** Modifying the state directly is not recommended and
 596        > could result in broken connections, and/or incorrect sync behavior.
 597
 598        Replaces the entire connection state with the provided state blob.
 599        Uses the safe variant that prevents updates while a sync is running (HTTP 423).
 600
 601        This is the counterpart to `dump_raw_state()` for backup/restore workflows.
 602        The `connectionId` in the blob is always overridden with this connection's
 603        ID, making state blobs portable across connections.
 604
 605        Accepts either format:
 606
 607        - **Config API format** (dict with `stateType`): passed through directly.
 608        - **Airbyte protocol format** (list of `AirbyteStateMessage` dicts): automatically
 609          converted to Config API format before sending.
 610
 611        Args:
 612            connection_state: Connection state in either Config API or Airbyte protocol format.
 613
 614        Returns:
 615            The updated connection state as a dictionary.
 616
 617        Raises:
 618            AirbyteConnectionSyncActiveError: If a sync is currently running on this
 619                connection (HTTP 423). Wait for the sync to complete before retrying.
 620        """
 621        api_state: dict[str, Any]
 622        if isinstance(connection_state, list):
 623            if not _is_protocol_state_format(connection_state):
 624                msg = (
 625                    "Expected connection_state list to contain Airbyte protocol state "
 626                    "message dicts (each with a top-level `type` of STREAM, GLOBAL, "
 627                    "or LEGACY). Got a list that does not match protocol format."
 628                )
 629                raise ValueError(msg)
 630            api_state = _denormalize_protocol_state_to_api(
 631                protocol_messages=connection_state,
 632                connection_id=self.connection_id,
 633            )
 634        elif isinstance(connection_state, dict):
 635            if _is_protocol_state_format(connection_state):
 636                api_state = _denormalize_protocol_state_to_api(
 637                    protocol_messages=[connection_state],
 638                    connection_id=self.connection_id,
 639                )
 640            else:
 641                api_state = connection_state
 642        else:
 643            msg = f"Expected a dict or list, got {type(connection_state)}"
 644            raise TypeError(msg)
 645
 646        return api_util.replace_connection_state(
 647            connection_id=self.connection_id,
 648            connection_state_dict=api_state,
 649            api_root=self.workspace.api_root,
 650            client_id=self.workspace.client_id,
 651            client_secret=self.workspace.client_secret,
 652            bearer_token=self.workspace.bearer_token,
 653            config_api_root=self.workspace.config_api_root,
 654        )
 655
 656    def get_stream_state(
 657        self,
 658        stream_name: str,
 659        stream_namespace: str | None = None,
 660    ) -> dict[str, Any] | None:
 661        """Get the state blob for a single stream within this connection.
 662
 663        Returns just the stream's state dictionary (e.g., {"cursor": "2024-01-01"}),
 664        not the full connection state envelope.
 665
 666        This is compatible with `stream`-type state and stream-level entries
 667        within a `global`-type state. It is not compatible with `legacy` state.
 668        To get or set the entire connection-level state artifact, use
 669        `dump_raw_state` and `import_raw_state` instead.
 670
 671        Args:
 672            stream_name: The name of the stream to get state for.
 673            stream_namespace: The source-side stream namespace. This refers to the
 674                namespace from the source (e.g., database schema), not any destination
 675                namespace override set in connection advanced settings.
 676
 677        Returns:
 678            The stream's state blob as a dictionary, or None if the stream is not found.
 679        """
 680        state_data = self.dump_raw_state(normalize=False)
 681        result = ConnectionStateResponse(**state_data)
 682
 683        streams = _get_stream_list(result)
 684        matching = [s for s in streams if _match_stream(s, stream_name, stream_namespace)]
 685
 686        if not matching:
 687            available = [s.stream_descriptor.name for s in streams]
 688            logger.warning(
 689                "Stream '%s' not found in connection state for connection '%s'. "
 690                "Available streams: %s",
 691                stream_name,
 692                self.connection_id,
 693                available,
 694            )
 695            return None
 696
 697        return matching[0].stream_state
 698
 699    def set_stream_state(
 700        self,
 701        stream_name: str,
 702        state_blob_dict: dict[str, Any],
 703        stream_namespace: str | None = None,
 704    ) -> None:
 705        """Set the state for a single stream within this connection.
 706
 707        Fetches the current full state, replaces only the specified stream's state,
 708        then sends the full updated state back to the API. If the stream does not
 709        exist in the current state, it is appended.
 710
 711        This is compatible with `stream`-type state and stream-level entries
 712        within a `global`-type state. It is not compatible with `legacy` state.
 713        To get or set the entire connection-level state artifact, use
 714        `dump_raw_state` and `import_raw_state` instead.
 715
 716        Uses the safe variant that prevents updates while a sync is running (HTTP 423).
 717
 718        Args:
 719            stream_name: The name of the stream to update state for.
 720            state_blob_dict: The state blob dict for this stream (e.g., {"cursor": "2024-01-01"}).
 721            stream_namespace: The source-side stream namespace. This refers to the
 722                namespace from the source (e.g., database schema), not any destination
 723                namespace override set in connection advanced settings.
 724
 725        Raises:
 726            PyAirbyteInputError: If the connection state type is not supported for
 727                stream-level operations (not_set, legacy).
 728            AirbyteConnectionSyncActiveError: If a sync is currently running on this
 729                connection (HTTP 423). Wait for the sync to complete before retrying.
 730        """
 731        state_data = self.dump_raw_state(normalize=False)
 732        current = ConnectionStateResponse(**state_data)
 733
 734        if current.state_type == "not_set":
 735            raise PyAirbyteInputError(
 736                message="Cannot set stream state: connection has no existing state.",
 737                context={"connection_id": self.connection_id},
 738            )
 739
 740        if current.state_type == "legacy":
 741            raise PyAirbyteInputError(
 742                message="Cannot set stream state on a legacy-type connection state.",
 743                context={"connection_id": self.connection_id},
 744            )
 745
 746        new_stream_entry = {
 747            "streamDescriptor": {
 748                "name": stream_name,
 749                **(
 750                    {
 751                        "namespace": stream_namespace,
 752                    }
 753                    if stream_namespace
 754                    else {}
 755                ),
 756            },
 757            "streamState": state_blob_dict,
 758        }
 759
 760        raw_streams: list[dict[str, Any]]
 761        if current.state_type == "stream":
 762            raw_streams = state_data.get("streamState", [])
 763        elif current.state_type == "global":
 764            raw_streams = state_data.get("globalState", {}).get("streamStates", [])
 765        else:
 766            raw_streams = []
 767
 768        streams = _get_stream_list(current)
 769        found = False
 770        updated_streams_raw: list[dict[str, Any]] = []
 771        for raw_s, parsed_s in zip(raw_streams, streams, strict=False):
 772            if _match_stream(parsed_s, stream_name, stream_namespace):
 773                updated_streams_raw.append(new_stream_entry)
 774                found = True
 775            else:
 776                updated_streams_raw.append(raw_s)
 777
 778        if not found:
 779            updated_streams_raw.append(new_stream_entry)
 780
 781        full_state: dict[str, Any] = {
 782            **state_data,
 783        }
 784
 785        if current.state_type == "stream":
 786            full_state["streamState"] = updated_streams_raw
 787        elif current.state_type == "global":
 788            original_global = state_data.get("globalState", {})
 789            full_state["globalState"] = {
 790                **original_global,
 791                "streamStates": updated_streams_raw,
 792            }
 793
 794        self.import_raw_state(full_state)
 795
 796    @deprecated("Use 'dump_raw_catalog()' instead.")
 797    def get_catalog_artifact(self) -> dict[str, Any] | None:
 798        """Get the configured catalog for this connection.
 799
 800        Returns the full configured catalog (syncCatalog) for this connection,
 801        including stream schemas, sync modes, cursor fields, and primary keys.
 802
 803        Uses the Config API endpoint: POST /v1/web_backend/connections/get
 804
 805        Returns:
 806            Dictionary containing the configured catalog, or `None` if not found.
 807        """
 808        return self.dump_raw_catalog()
 809
 810    def dump_raw_catalog(
 811        self,
 812        *,
 813        normalize: bool = True,
 814    ) -> dict[str, Any] | None:
 815        """Dump the configured catalog for this connection.
 816
 817        By default, returns the catalog in Airbyte protocol format
 818        (`ConfiguredAirbyteCatalog` with snake_case keys), suitable for passing
 819        to a connector's `--catalog` flag.
 820
 821        When `normalize` is `False`, returns the raw `syncCatalog` dict from the
 822        Config API (camelCase keys, nested `config` block). This raw format can be
 823        passed directly to `import_raw_catalog()` for backup/restore workflows.
 824
 825        Args:
 826            normalize: If `True` (default), convert to Airbyte protocol format.
 827                If `False`, return the raw Config API catalog.
 828
 829        Returns:
 830            The configured catalog dict, or `None` if not found.
 831        """
 832        connection_response = api_util.get_connection_catalog(
 833            connection_id=self.connection_id,
 834            api_root=self.workspace.api_root,
 835            client_id=self.workspace.client_id,
 836            client_secret=self.workspace.client_secret,
 837            bearer_token=self.workspace.bearer_token,
 838            config_api_root=self.workspace.config_api_root,
 839        )
 840        raw = connection_response.get("syncCatalog")
 841        if raw is None:
 842            return None
 843        if normalize:
 844            return _normalize_catalog_to_protocol(raw)
 845        return raw
 846
 847    def import_raw_catalog(self, catalog: dict[str, Any]) -> None:
 848        """Replace the configured catalog for this connection.
 849
 850        > ⚠️ **WARNING:** Modifying the catalog directly is not recommended and
 851        > could result in broken connections, and/or incorrect sync behavior.
 852
 853        Accepts a configured catalog dict and replaces the connection's entire
 854        catalog with it. All other connection settings remain unchanged.
 855
 856        Accepts either format:
 857
 858        - **Config API format** (`syncCatalog` with camelCase keys and nested `config`):
 859          passed through directly.
 860        - **Airbyte protocol format** (`ConfiguredAirbyteCatalog` with snake_case keys):
 861          automatically converted to Config API format before sending.
 862
 863        Args:
 864            catalog: The configured catalog dict in either format.
 865        """
 866        if _is_protocol_catalog_format(catalog):
 867            catalog = _denormalize_catalog_to_api(catalog)
 868
 869        api_util.replace_connection_catalog(
 870            connection_id=self.connection_id,
 871            configured_catalog_dict=catalog,
 872            api_root=self.workspace.api_root,
 873            client_id=self.workspace.client_id,
 874            client_secret=self.workspace.client_secret,
 875            bearer_token=self.workspace.bearer_token,
 876            config_api_root=self.workspace.config_api_root,
 877        )
 878
 879    def rename(self, name: str) -> CloudConnection:
 880        """Rename the connection.
 881
 882        Args:
 883            name: New name for the connection
 884
 885        Returns:
 886            Updated CloudConnection object with refreshed info
 887        """
 888        updated_response = api_util.patch_connection(
 889            connection_id=self.connection_id,
 890            api_root=self.workspace.api_root,
 891            client_id=self.workspace.client_id,
 892            client_secret=self.workspace.client_secret,
 893            bearer_token=self.workspace.bearer_token,
 894            name=name,
 895        )
 896        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
 897        return self
 898
 899    def set_table_prefix(self, prefix: str) -> CloudConnection:
 900        """Set the table prefix for the connection.
 901
 902        Args:
 903            prefix: New table prefix to use when syncing to the destination
 904
 905        Returns:
 906            Updated CloudConnection object with refreshed info
 907        """
 908        updated_response = api_util.patch_connection(
 909            connection_id=self.connection_id,
 910            api_root=self.workspace.api_root,
 911            client_id=self.workspace.client_id,
 912            client_secret=self.workspace.client_secret,
 913            bearer_token=self.workspace.bearer_token,
 914            prefix=prefix,
 915        )
 916        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
 917        return self
 918
 919    def set_selected_streams(self, stream_names: list[str]) -> CloudConnection:
 920        """Set the selected streams for the connection.
 921
 922        This is a destructive operation that can break existing connections if the
 923        stream selection is changed incorrectly. Use with caution.
 924
 925        Args:
 926            stream_names: List of stream names to sync
 927
 928        Returns:
 929            Updated CloudConnection object with refreshed info
 930        """
 931        configurations = api_util.build_stream_configurations(stream_names)
 932
 933        updated_response = api_util.patch_connection(
 934            connection_id=self.connection_id,
 935            api_root=self.workspace.api_root,
 936            client_id=self.workspace.client_id,
 937            client_secret=self.workspace.client_secret,
 938            bearer_token=self.workspace.bearer_token,
 939            configurations=configurations,
 940        )
 941        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
 942        return self
 943
 944    # Enable/Disable
 945
 946    @property
 947    def enabled(self) -> bool:
 948        """Get the current enabled status of the connection.
 949
 950        This property always fetches fresh data from the API to ensure accuracy,
 951        as another process or user may have toggled the setting.
 952
 953        Returns:
 954            True if the connection status is 'active', False otherwise.
 955        """
 956        connection_info = self._fetch_connection_info(force_refresh=True)
 957        return connection_info.status == "active"
 958
 959    @enabled.setter
 960    def enabled(self, value: bool) -> None:
 961        """Set the enabled status of the connection.
 962
 963        Args:
 964            value: True to enable (set status to 'active'), False to disable
 965                (set status to 'inactive').
 966        """
 967        self.set_enabled(enabled=value)
 968
 969    def set_enabled(
 970        self,
 971        *,
 972        enabled: bool,
 973        ignore_noop: bool = True,
 974    ) -> None:
 975        """Set the enabled status of the connection.
 976
 977        Args:
 978            enabled: True to enable (set status to 'active'), False to disable
 979                (set status to 'inactive').
 980            ignore_noop: If True (default), silently return if the connection is already
 981                in the requested state. If False, raise ValueError when the requested
 982                state matches the current state.
 983
 984        Raises:
 985            ValueError: If ignore_noop is False and the connection is already in the
 986                requested state.
 987        """
 988        # Always fetch fresh data to check current status
 989        connection_info = self._fetch_connection_info(force_refresh=True)
 990        current_status = connection_info.status
 991        desired_status = "active" if enabled else "inactive"
 992
 993        if current_status == desired_status:
 994            if ignore_noop:
 995                return
 996            raise ValueError(
 997                f"Connection is already {'enabled' if enabled else 'disabled'}. "
 998                f"Current status: {current_status}"
 999            )
1000
1001        updated_response = api_util.patch_connection(
1002            connection_id=self.connection_id,
1003            api_root=self.workspace.api_root,
1004            client_id=self.workspace.client_id,
1005            client_secret=self.workspace.client_secret,
1006            bearer_token=self.workspace.bearer_token,
1007            status=desired_status,
1008        )
1009        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
1010
1011    # Scheduling
1012
1013    def set_schedule(
1014        self,
1015        cron_expression: str,
1016    ) -> None:
1017        """Set a cron schedule for the connection.
1018
1019        Args:
1020            cron_expression: A Quartz cron expression defining when syncs should run.
1021                Quartz expressions have 6 or 7 space-separated fields
1022                (seconds, minutes, hours, day-of-month, month, day-of-week[, year]),
1023                optionally followed by a timezone ID. The Airbyte API rejects standard
1024                5-field Unix cron expressions and schedules that run more often than
1025                once per hour.
1026
1027        Examples:
1028            - "0 0 0 * * ?"  # Daily at midnight UTC
1029            - "0 0 */6 * * ?"  # Every 6 hours
1030            - "0 0 0 ? * SUN"  # Weekly on Sunday at midnight UTC
1031            - "0 0 9 ? * MON-FRI US/Pacific"  # Weekdays at 9am Pacific
1032        """
1033        _validate_quartz_cron_expression(cron_expression)
1034        updated_response = api_util.patch_connection(
1035            connection_id=self.connection_id,
1036            api_root=self.workspace.api_root,
1037            client_id=self.workspace.client_id,
1038            client_secret=self.workspace.client_secret,
1039            bearer_token=self.workspace.bearer_token,
1040            schedule=api_util.build_connection_schedule(
1041                schedule_type="cron",
1042                cron_expression=cron_expression,
1043            ),
1044        )
1045        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
1046
1047    def set_manual_schedule(self) -> None:
1048        """Set the connection to manual scheduling.
1049
1050        Disables automatic syncs. Syncs will only run when manually triggered.
1051        """
1052        updated_response = api_util.patch_connection(
1053            connection_id=self.connection_id,
1054            api_root=self.workspace.api_root,
1055            client_id=self.workspace.client_id,
1056            client_secret=self.workspace.client_secret,
1057            bearer_token=self.workspace.bearer_token,
1058            schedule=api_util.build_connection_schedule(schedule_type="manual"),
1059        )
1060        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
1061
1062    # Deletions
1063
1064    def permanently_delete(
1065        self,
1066        *,
1067        cascade_delete_source: bool = False,
1068        cascade_delete_destination: bool = False,
1069    ) -> None:
1070        """Delete the connection.
1071
1072        Args:
1073            cascade_delete_source: Whether to also delete the source.
1074            cascade_delete_destination: Whether to also delete the destination.
1075        """
1076        self.workspace.permanently_delete_connection(self)
1077
1078        if cascade_delete_source:
1079            self.workspace.permanently_delete_source(self.source_id)
1080
1081        if cascade_delete_destination:
1082            self.workspace.permanently_delete_destination(self.destination_id)

A connection is an extract-load (EL) pairing of a source and destination in Airbyte Cloud.

You can use a connection object to run sync jobs, retrieve logs, and manage the connection.

CloudConnection( workspace: CloudWorkspace, connection_id: str, source: str | None = None, destination: str | None = None)
 83    def __init__(
 84        self,
 85        workspace: CloudWorkspace,
 86        connection_id: str,
 87        source: str | None = None,
 88        destination: str | None = None,
 89    ) -> None:
 90        """It is not recommended to create a `CloudConnection` object directly.
 91
 92        Instead, use `CloudWorkspace.get_connection()` to create a connection object.
 93        """
 94        self.connection_id = connection_id
 95        """The ID of the connection."""
 96
 97        self.workspace = workspace
 98        """The workspace that the connection belongs to."""
 99
100        self._source_id = source
101        """The ID of the source."""
102
103        self._destination_id = destination
104        """The ID of the destination."""
105
106        self._connection_info: CloudConnectionInfo | None = None
107        """The connection info object. (Cached.)"""
108
109        self._cloud_source_object: CloudSource | None = None
110        """The source object. (Cached.)"""
111
112        self._cloud_destination_object: CloudDestination | None = None
113        """The destination object. (Cached.)"""

It is not recommended to create a CloudConnection object directly.

Instead, use CloudWorkspace.get_connection() to create a connection object.

connection_id

The ID of the connection.

workspace

The workspace that the connection belongs to.

def check_is_valid(self) -> bool:
186    def check_is_valid(self) -> bool:
187        """Check if this connection exists and belongs to the expected workspace.
188
189        This method fetches connection info from the API (if not already cached) and
190        verifies that the connection's workspace_id matches the workspace associated
191        with this CloudConnection object.
192
193        Returns:
194            True if the connection exists and belongs to the expected workspace.
195
196        Raises:
197            AirbyteWorkspaceMismatchError: If the connection belongs to a different workspace.
198            AirbyteMissingResourceError: If the connection doesn't exist.
199        """
200        self._fetch_connection_info(force_refresh=False, verify=True)
201        return True

Check if this connection exists and belongs to the expected workspace.

This method fetches connection info from the API (if not already cached) and verifies that the connection's workspace_id matches the workspace associated with this CloudConnection object.

Returns:

True if the connection exists and belongs to the expected workspace.

Raises:
  • AirbyteWorkspaceMismatchError: If the connection belongs to a different workspace.
  • AirbyteMissingResourceError: If the connection doesn't exist.
name: str | None
222    @property
223    def name(self) -> str | None:
224        """Get the display name of the connection, if available.
225
226        E.g. "My Postgres to Snowflake", not the connection ID.
227        """
228        if not self._connection_info:
229            self._connection_info = self._fetch_connection_info()
230
231        return self._connection_info.name

Get the display name of the connection, if available.

E.g. "My Postgres to Snowflake", not the connection ID.

source_id: str
233    @property
234    def source_id(self) -> str:
235        """The ID of the source."""
236        if not self._source_id:
237            if not self._connection_info:
238                self._connection_info = self._fetch_connection_info()
239
240            self._source_id = self._connection_info.source_id
241
242        return self._source_id

The ID of the source.

source: airbyte.cloud.connectors.CloudSource
244    @property
245    def source(self) -> CloudSource:
246        """Get the source object."""
247        if self._cloud_source_object:
248            return self._cloud_source_object
249
250        self._cloud_source_object = CloudSource(
251            workspace=self.workspace,
252            connector_id=self.source_id,
253        )
254        return self._cloud_source_object

Get the source object.

destination_id: str
256    @property
257    def destination_id(self) -> str:
258        """The ID of the destination."""
259        if not self._destination_id:
260            if not self._connection_info:
261                self._connection_info = self._fetch_connection_info()
262
263            self._destination_id = self._connection_info.destination_id
264
265        return self._destination_id

The ID of the destination.

destination: airbyte.cloud.connectors.CloudDestination
267    @property
268    def destination(self) -> CloudDestination:
269        """Get the destination object."""
270        if self._cloud_destination_object:
271            return self._cloud_destination_object
272
273        self._cloud_destination_object = CloudDestination(
274            workspace=self.workspace,
275            connector_id=self.destination_id,
276        )
277        return self._cloud_destination_object

Get the destination object.

stream_names: list[str]
279    @property
280    def stream_names(self) -> list[str]:
281        """The stream names."""
282        if not self._connection_info:
283            self._connection_info = self._fetch_connection_info()
284
285        return [stream.name for stream in self._connection_info.configurations.streams or []]

The stream names.

table_prefix: str
287    @property
288    def table_prefix(self) -> str:
289        """The table prefix."""
290        if not self._connection_info:
291            self._connection_info = self._fetch_connection_info()
292
293        return self._connection_info.prefix or ""

The table prefix.

namespace_definition: str | None
295    @property
296    def namespace_definition(self) -> str | None:
297        """How destination namespaces are chosen: `source`, `destination`, or `custom_format`."""
298        if not self._connection_info:
299            self._connection_info = self._fetch_connection_info()
300
301        return self._connection_info.namespace_definition

How destination namespaces are chosen: source, destination, or custom_format.

namespace_format: str | None
303    @property
304    def namespace_format(self) -> str | None:
305        """The namespace format template, when `namespace_definition` is `custom_format`."""
306        if not self._connection_info:
307            self._connection_info = self._fetch_connection_info()
308
309        return self._connection_info.namespace_format

The namespace format template, when namespace_definition is custom_format.

connection_url: str | None
311    @property
312    def connection_url(self) -> str | None:
313        """The web URL to the connection."""
314        return f"{self.workspace.workspace_url}/connections/{self.connection_id}"

The web URL to the connection.

job_history_url: str | None
316    @property
317    def job_history_url(self) -> str | None:
318        """The URL to the job history for the connection."""
319        return f"{self.connection_url}/timeline"

The URL to the job history for the connection.

def run_sync( self, *, wait: bool = True, wait_timeout: int = 300) -> SyncResult:
323    def run_sync(
324        self,
325        *,
326        wait: bool = True,
327        wait_timeout: int = 300,
328    ) -> SyncResult:
329        """Run a sync."""
330        try:
331            connection_response = api_util.run_connection(
332                connection_id=self.connection_id,
333                api_root=self.workspace.api_root,
334                workspace_id=self.workspace.workspace_id,
335                client_id=self.workspace.client_id,
336                client_secret=self.workspace.client_secret,
337                bearer_token=self.workspace.bearer_token,
338            )
339        except AirbyteConnectionSyncError as ex:
340            if (
341                ex.context
342                and ex.context.get("status_code") == HTTPStatus.CONFLICT
343                and not self.enabled
344            ):
345                raise PyAirbyteInputError(
346                    message=(
347                        f"Connection '{self.connection_id}' is disabled (status 'inactive'), "
348                        "so a sync cannot be started."
349                    ),
350                    guidance=(
351                        "Re-enable the connection first (e.g. "
352                        "`connection.set_enabled(enabled=True)`, or the "
353                        "`update_cloud_connection` MCP tool with `enabled=True`), then retry."
354                    ),
355                    context={"connection_id": self.connection_id},
356                ) from ex
357            raise
358        sync_result = SyncResult(
359            workspace=self.workspace,
360            connection=self,
361            job_id=connection_response.job_id,
362        )
363
364        if wait:
365            sync_result.wait_for_completion(
366                wait_timeout=wait_timeout,
367                raise_failure=True,
368                raise_timeout=True,
369            )
370
371        return sync_result

Run a sync.

def cancel_sync(self, job_id: int | None = None) -> SyncResult:
417    def cancel_sync(self, job_id: int | None = None) -> SyncResult:
418        """Cancel a running sync job.
419
420        Defaults to the connection's most recent sync job. Other job types must be
421        targeted with an explicit `job_id`.
422        """
423        target_job_id: int = (
424            self._get_latest_cancellable_sync_job_id()
425            if job_id is None
426            else self._validated_cancellable_job_id(job_id)
427        )
428
429        job_response = api_util.cancel_job(
430            job_id=target_job_id,
431            api_root=self.workspace.api_root,
432            client_id=self.workspace.client_id,
433            client_secret=self.workspace.client_secret,
434            bearer_token=self.workspace.bearer_token,
435        )
436        return SyncResult(
437            workspace=self.workspace,
438            connection=self,
439            job_id=job_response.job_id,
440            _latest_job_info=CloudJobInfo.from_api_response(job_response),
441        )

Cancel a running sync job.

Defaults to the connection's most recent sync job. Other job types must be targeted with an explicit job_id.

def get_previous_sync_logs( self, *, limit: int = 20, offset: int | None = None, from_tail: bool = True, job_type: str | JobTypeEnum | None = None) -> list[SyncResult]:
452    def get_previous_sync_logs(
453        self,
454        *,
455        limit: int = 20,
456        offset: int | None = None,
457        from_tail: bool = True,
458        job_type: str | JobTypeEnum | None = None,
459    ) -> list[SyncResult]:
460        """Get previous sync jobs for a connection with pagination support.
461
462        Returns SyncResult objects containing job metadata (job_id, status, bytes_synced,
463        rows_synced, start_time). Full log text can be fetched lazily via
464        `SyncResult.get_full_log_text()`.
465
466        Args:
467            limit: Maximum number of jobs to return. Defaults to 20.
468            offset: Number of jobs to skip from the beginning. Defaults to None (0).
469            from_tail: If True, returns jobs ordered newest-first (createdAt DESC).
470                If False, returns jobs ordered oldest-first (createdAt ASC).
471                Defaults to True.
472            job_type: Filter by job type (e.g., `sync`, `refresh`).
473                If not specified, defaults to sync and reset jobs only (API default behavior).
474
475        Returns:
476            A list of SyncResult objects representing the sync jobs.
477        """
478        order_by = (
479            api_util.JOB_ORDER_BY_CREATED_AT_DESC
480            if from_tail
481            else api_util.JOB_ORDER_BY_CREATED_AT_ASC
482        )
483        sync_logs = api_util.get_job_logs(
484            connection_id=self.connection_id,
485            api_root=self.workspace.api_root,
486            workspace_id=self.workspace.workspace_id,
487            limit=limit,
488            offset=offset,
489            order_by=order_by,
490            job_type=job_type,
491            client_id=self.workspace.client_id,
492            client_secret=self.workspace.client_secret,
493            bearer_token=self.workspace.bearer_token,
494        )
495        return [
496            SyncResult(
497                workspace=self.workspace,
498                connection=self,
499                job_id=sync_log.job_id,
500                _latest_job_info=CloudJobInfo.from_api_response(sync_log),
501            )
502            for sync_log in sync_logs
503        ]

Get previous sync jobs for a connection with pagination support.

Returns SyncResult objects containing job metadata (job_id, status, bytes_synced, rows_synced, start_time). Full log text can be fetched lazily via SyncResult.get_full_log_text().

Arguments:
  • limit: Maximum number of jobs to return. Defaults to 20.
  • offset: Number of jobs to skip from the beginning. Defaults to None (0).
  • from_tail: If True, returns jobs ordered newest-first (createdAt DESC). If False, returns jobs ordered oldest-first (createdAt ASC). Defaults to True.
  • job_type: Filter by job type (e.g., sync, refresh). If not specified, defaults to sync and reset jobs only (API default behavior).
Returns:

A list of SyncResult objects representing the sync jobs.

def get_sync_result( self, job_id: int | None = None) -> SyncResult | None:
505    def get_sync_result(
506        self,
507        job_id: int | None = None,
508    ) -> SyncResult | None:
509        """Get the sync result for the connection.
510
511        If `job_id` is not provided, the most recent sync job will be used.
512
513        Returns `None` if job_id is omitted and no previous jobs are found.
514        """
515        if job_id is None:
516            # Get the most recent sync job
517            results = self.get_previous_sync_logs(
518                limit=1,
519            )
520            if results:
521                return results[0]
522
523            return None
524
525        # Get the sync job by ID (lazy loaded)
526        return SyncResult(
527            workspace=self.workspace,
528            connection=self,
529            job_id=job_id,
530        )

Get the sync result for the connection.

If job_id is not provided, the most recent sync job will be used.

Returns None if job_id is omitted and no previous jobs are found.

@deprecated("Use 'dump_raw_state()' instead.")
def get_state_artifacts(self) -> list[dict[str, typing.Any]] | None:
534    @deprecated("Use 'dump_raw_state()' instead.")
535    def get_state_artifacts(self) -> list[dict[str, Any]] | None:
536        """Deprecated. Use `dump_raw_state()` instead."""
537        state_response = api_util.get_connection_state(
538            connection_id=self.connection_id,
539            api_root=self.workspace.api_root,
540            client_id=self.workspace.client_id,
541            client_secret=self.workspace.client_secret,
542            bearer_token=self.workspace.bearer_token,
543            config_api_root=self.workspace.config_api_root,
544        )
545        if state_response.get("stateType") == "not_set":
546            return None
547        return state_response.get("streamState", [])

Deprecated. Use dump_raw_state() instead.

def dump_raw_state( self, *, normalize: bool = True) -> dict[str, typing.Any] | list[dict[str, typing.Any]]:
555    def dump_raw_state(
556        self,
557        *,
558        normalize: bool = True,
559    ) -> dict[str, Any] | list[dict[str, Any]]:
560        """Dump the state for this connection.
561
562        By default, returns a list of Airbyte protocol `AirbyteStateMessage` dicts
563        with snake_case keys, suitable for passing to a connector's `--state` flag.
564
565        When `normalize` is `False`, returns the raw Config API dict (camelCase keys,
566        includes `stateType` and `connectionId`). This raw format can be passed
567        directly to `import_raw_state()` for backup/restore workflows.
568
569        Args:
570            normalize: If `True` (default), convert to Airbyte protocol format.
571                If `False`, return the raw Config API response.
572
573        Returns:
574            Normalized: list of protocol-format state message dicts (empty list if
575            no state). Raw: the full Config API state dict.
576        """
577        raw = api_util.get_connection_state(
578            connection_id=self.connection_id,
579            api_root=self.workspace.api_root,
580            client_id=self.workspace.client_id,
581            client_secret=self.workspace.client_secret,
582            bearer_token=self.workspace.bearer_token,
583            config_api_root=self.workspace.config_api_root,
584        )
585        if normalize:
586            return _normalize_state_to_protocol(raw)
587        return raw

Dump the state for this connection.

By default, returns a list of Airbyte protocol AirbyteStateMessage dicts with snake_case keys, suitable for passing to a connector's --state flag.

When normalize is False, returns the raw Config API dict (camelCase keys, includes stateType and connectionId). This raw format can be passed directly to import_raw_state() for backup/restore workflows.

Arguments:
  • normalize: If True (default), convert to Airbyte protocol format. If False, return the raw Config API response.
Returns:

Normalized: list of protocol-format state message dicts (empty list if no state). Raw: the full Config API state dict.

def import_raw_state( self, connection_state: dict[str, typing.Any] | list[dict[str, typing.Any]]) -> dict[str, typing.Any]:
589    def import_raw_state(
590        self,
591        connection_state: dict[str, Any] | list[dict[str, Any]],
592    ) -> dict[str, Any]:
593        """Import (restore) the full state for this connection.
594
595        > ⚠️ **WARNING:** Modifying the state directly is not recommended and
596        > could result in broken connections, and/or incorrect sync behavior.
597
598        Replaces the entire connection state with the provided state blob.
599        Uses the safe variant that prevents updates while a sync is running (HTTP 423).
600
601        This is the counterpart to `dump_raw_state()` for backup/restore workflows.
602        The `connectionId` in the blob is always overridden with this connection's
603        ID, making state blobs portable across connections.
604
605        Accepts either format:
606
607        - **Config API format** (dict with `stateType`): passed through directly.
608        - **Airbyte protocol format** (list of `AirbyteStateMessage` dicts): automatically
609          converted to Config API format before sending.
610
611        Args:
612            connection_state: Connection state in either Config API or Airbyte protocol format.
613
614        Returns:
615            The updated connection state as a dictionary.
616
617        Raises:
618            AirbyteConnectionSyncActiveError: If a sync is currently running on this
619                connection (HTTP 423). Wait for the sync to complete before retrying.
620        """
621        api_state: dict[str, Any]
622        if isinstance(connection_state, list):
623            if not _is_protocol_state_format(connection_state):
624                msg = (
625                    "Expected connection_state list to contain Airbyte protocol state "
626                    "message dicts (each with a top-level `type` of STREAM, GLOBAL, "
627                    "or LEGACY). Got a list that does not match protocol format."
628                )
629                raise ValueError(msg)
630            api_state = _denormalize_protocol_state_to_api(
631                protocol_messages=connection_state,
632                connection_id=self.connection_id,
633            )
634        elif isinstance(connection_state, dict):
635            if _is_protocol_state_format(connection_state):
636                api_state = _denormalize_protocol_state_to_api(
637                    protocol_messages=[connection_state],
638                    connection_id=self.connection_id,
639                )
640            else:
641                api_state = connection_state
642        else:
643            msg = f"Expected a dict or list, got {type(connection_state)}"
644            raise TypeError(msg)
645
646        return api_util.replace_connection_state(
647            connection_id=self.connection_id,
648            connection_state_dict=api_state,
649            api_root=self.workspace.api_root,
650            client_id=self.workspace.client_id,
651            client_secret=self.workspace.client_secret,
652            bearer_token=self.workspace.bearer_token,
653            config_api_root=self.workspace.config_api_root,
654        )

Import (restore) the full state for this connection.

⚠️ WARNING: Modifying the state directly is not recommended and could result in broken connections, and/or incorrect sync behavior.

Replaces the entire connection state with the provided state blob. Uses the safe variant that prevents updates while a sync is running (HTTP 423).

This is the counterpart to dump_raw_state() for backup/restore workflows. The connectionId in the blob is always overridden with this connection's ID, making state blobs portable across connections.

Accepts either format:

  • Config API format (dict with stateType): passed through directly.
  • Airbyte protocol format (list of AirbyteStateMessage dicts): automatically converted to Config API format before sending.
Arguments:
  • connection_state: Connection state in either Config API or Airbyte protocol format.
Returns:

The updated connection state as a dictionary.

Raises:
  • AirbyteConnectionSyncActiveError: If a sync is currently running on this connection (HTTP 423). Wait for the sync to complete before retrying.
def get_stream_state( self, stream_name: str, stream_namespace: str | None = None) -> dict[str, typing.Any] | None:
656    def get_stream_state(
657        self,
658        stream_name: str,
659        stream_namespace: str | None = None,
660    ) -> dict[str, Any] | None:
661        """Get the state blob for a single stream within this connection.
662
663        Returns just the stream's state dictionary (e.g., {"cursor": "2024-01-01"}),
664        not the full connection state envelope.
665
666        This is compatible with `stream`-type state and stream-level entries
667        within a `global`-type state. It is not compatible with `legacy` state.
668        To get or set the entire connection-level state artifact, use
669        `dump_raw_state` and `import_raw_state` instead.
670
671        Args:
672            stream_name: The name of the stream to get state for.
673            stream_namespace: The source-side stream namespace. This refers to the
674                namespace from the source (e.g., database schema), not any destination
675                namespace override set in connection advanced settings.
676
677        Returns:
678            The stream's state blob as a dictionary, or None if the stream is not found.
679        """
680        state_data = self.dump_raw_state(normalize=False)
681        result = ConnectionStateResponse(**state_data)
682
683        streams = _get_stream_list(result)
684        matching = [s for s in streams if _match_stream(s, stream_name, stream_namespace)]
685
686        if not matching:
687            available = [s.stream_descriptor.name for s in streams]
688            logger.warning(
689                "Stream '%s' not found in connection state for connection '%s'. "
690                "Available streams: %s",
691                stream_name,
692                self.connection_id,
693                available,
694            )
695            return None
696
697        return matching[0].stream_state

Get the state blob for a single stream within this connection.

Returns just the stream's state dictionary (e.g., {"cursor": "2024-01-01"}), not the full connection state envelope.

This is compatible with stream-type state and stream-level entries within a global-type state. It is not compatible with legacy state. To get or set the entire connection-level state artifact, use dump_raw_state and import_raw_state instead.

Arguments:
  • stream_name: The name of the stream to get state for.
  • stream_namespace: The source-side stream namespace. This refers to the namespace from the source (e.g., database schema), not any destination namespace override set in connection advanced settings.
Returns:

The stream's state blob as a dictionary, or None if the stream is not found.

def set_stream_state( self, stream_name: str, state_blob_dict: dict[str, typing.Any], stream_namespace: str | None = None) -> None:
699    def set_stream_state(
700        self,
701        stream_name: str,
702        state_blob_dict: dict[str, Any],
703        stream_namespace: str | None = None,
704    ) -> None:
705        """Set the state for a single stream within this connection.
706
707        Fetches the current full state, replaces only the specified stream's state,
708        then sends the full updated state back to the API. If the stream does not
709        exist in the current state, it is appended.
710
711        This is compatible with `stream`-type state and stream-level entries
712        within a `global`-type state. It is not compatible with `legacy` state.
713        To get or set the entire connection-level state artifact, use
714        `dump_raw_state` and `import_raw_state` instead.
715
716        Uses the safe variant that prevents updates while a sync is running (HTTP 423).
717
718        Args:
719            stream_name: The name of the stream to update state for.
720            state_blob_dict: The state blob dict for this stream (e.g., {"cursor": "2024-01-01"}).
721            stream_namespace: The source-side stream namespace. This refers to the
722                namespace from the source (e.g., database schema), not any destination
723                namespace override set in connection advanced settings.
724
725        Raises:
726            PyAirbyteInputError: If the connection state type is not supported for
727                stream-level operations (not_set, legacy).
728            AirbyteConnectionSyncActiveError: If a sync is currently running on this
729                connection (HTTP 423). Wait for the sync to complete before retrying.
730        """
731        state_data = self.dump_raw_state(normalize=False)
732        current = ConnectionStateResponse(**state_data)
733
734        if current.state_type == "not_set":
735            raise PyAirbyteInputError(
736                message="Cannot set stream state: connection has no existing state.",
737                context={"connection_id": self.connection_id},
738            )
739
740        if current.state_type == "legacy":
741            raise PyAirbyteInputError(
742                message="Cannot set stream state on a legacy-type connection state.",
743                context={"connection_id": self.connection_id},
744            )
745
746        new_stream_entry = {
747            "streamDescriptor": {
748                "name": stream_name,
749                **(
750                    {
751                        "namespace": stream_namespace,
752                    }
753                    if stream_namespace
754                    else {}
755                ),
756            },
757            "streamState": state_blob_dict,
758        }
759
760        raw_streams: list[dict[str, Any]]
761        if current.state_type == "stream":
762            raw_streams = state_data.get("streamState", [])
763        elif current.state_type == "global":
764            raw_streams = state_data.get("globalState", {}).get("streamStates", [])
765        else:
766            raw_streams = []
767
768        streams = _get_stream_list(current)
769        found = False
770        updated_streams_raw: list[dict[str, Any]] = []
771        for raw_s, parsed_s in zip(raw_streams, streams, strict=False):
772            if _match_stream(parsed_s, stream_name, stream_namespace):
773                updated_streams_raw.append(new_stream_entry)
774                found = True
775            else:
776                updated_streams_raw.append(raw_s)
777
778        if not found:
779            updated_streams_raw.append(new_stream_entry)
780
781        full_state: dict[str, Any] = {
782            **state_data,
783        }
784
785        if current.state_type == "stream":
786            full_state["streamState"] = updated_streams_raw
787        elif current.state_type == "global":
788            original_global = state_data.get("globalState", {})
789            full_state["globalState"] = {
790                **original_global,
791                "streamStates": updated_streams_raw,
792            }
793
794        self.import_raw_state(full_state)

Set the state for a single stream within this connection.

Fetches the current full state, replaces only the specified stream's state, then sends the full updated state back to the API. If the stream does not exist in the current state, it is appended.

This is compatible with stream-type state and stream-level entries within a global-type state. It is not compatible with legacy state. To get or set the entire connection-level state artifact, use dump_raw_state and import_raw_state instead.

Uses the safe variant that prevents updates while a sync is running (HTTP 423).

Arguments:
  • stream_name: The name of the stream to update state for.
  • state_blob_dict: The state blob dict for this stream (e.g., {"cursor": "2024-01-01"}).
  • stream_namespace: The source-side stream namespace. This refers to the namespace from the source (e.g., database schema), not any destination namespace override set in connection advanced settings.
Raises:
  • PyAirbyteInputError: If the connection state type is not supported for stream-level operations (not_set, legacy).
  • AirbyteConnectionSyncActiveError: If a sync is currently running on this connection (HTTP 423). Wait for the sync to complete before retrying.
@deprecated("Use 'dump_raw_catalog()' instead.")
def get_catalog_artifact(self) -> dict[str, typing.Any] | None:
796    @deprecated("Use 'dump_raw_catalog()' instead.")
797    def get_catalog_artifact(self) -> dict[str, Any] | None:
798        """Get the configured catalog for this connection.
799
800        Returns the full configured catalog (syncCatalog) for this connection,
801        including stream schemas, sync modes, cursor fields, and primary keys.
802
803        Uses the Config API endpoint: POST /v1/web_backend/connections/get
804
805        Returns:
806            Dictionary containing the configured catalog, or `None` if not found.
807        """
808        return self.dump_raw_catalog()

Get the configured catalog for this connection.

Returns the full configured catalog (syncCatalog) for this connection, including stream schemas, sync modes, cursor fields, and primary keys.

Uses the Config API endpoint: POST /v1/web_backend/connections/get

Returns:

Dictionary containing the configured catalog, or None if not found.

def dump_raw_catalog(self, *, normalize: bool = True) -> dict[str, typing.Any] | None:
810    def dump_raw_catalog(
811        self,
812        *,
813        normalize: bool = True,
814    ) -> dict[str, Any] | None:
815        """Dump the configured catalog for this connection.
816
817        By default, returns the catalog in Airbyte protocol format
818        (`ConfiguredAirbyteCatalog` with snake_case keys), suitable for passing
819        to a connector's `--catalog` flag.
820
821        When `normalize` is `False`, returns the raw `syncCatalog` dict from the
822        Config API (camelCase keys, nested `config` block). This raw format can be
823        passed directly to `import_raw_catalog()` for backup/restore workflows.
824
825        Args:
826            normalize: If `True` (default), convert to Airbyte protocol format.
827                If `False`, return the raw Config API catalog.
828
829        Returns:
830            The configured catalog dict, or `None` if not found.
831        """
832        connection_response = api_util.get_connection_catalog(
833            connection_id=self.connection_id,
834            api_root=self.workspace.api_root,
835            client_id=self.workspace.client_id,
836            client_secret=self.workspace.client_secret,
837            bearer_token=self.workspace.bearer_token,
838            config_api_root=self.workspace.config_api_root,
839        )
840        raw = connection_response.get("syncCatalog")
841        if raw is None:
842            return None
843        if normalize:
844            return _normalize_catalog_to_protocol(raw)
845        return raw

Dump the configured catalog for this connection.

By default, returns the catalog in Airbyte protocol format (ConfiguredAirbyteCatalog with snake_case keys), suitable for passing to a connector's --catalog flag.

When normalize is False, returns the raw syncCatalog dict from the Config API (camelCase keys, nested config block). This raw format can be passed directly to import_raw_catalog() for backup/restore workflows.

Arguments:
  • normalize: If True (default), convert to Airbyte protocol format. If False, return the raw Config API catalog.
Returns:

The configured catalog dict, or None if not found.

def import_raw_catalog(self, catalog: dict[str, typing.Any]) -> None:
847    def import_raw_catalog(self, catalog: dict[str, Any]) -> None:
848        """Replace the configured catalog for this connection.
849
850        > ⚠️ **WARNING:** Modifying the catalog directly is not recommended and
851        > could result in broken connections, and/or incorrect sync behavior.
852
853        Accepts a configured catalog dict and replaces the connection's entire
854        catalog with it. All other connection settings remain unchanged.
855
856        Accepts either format:
857
858        - **Config API format** (`syncCatalog` with camelCase keys and nested `config`):
859          passed through directly.
860        - **Airbyte protocol format** (`ConfiguredAirbyteCatalog` with snake_case keys):
861          automatically converted to Config API format before sending.
862
863        Args:
864            catalog: The configured catalog dict in either format.
865        """
866        if _is_protocol_catalog_format(catalog):
867            catalog = _denormalize_catalog_to_api(catalog)
868
869        api_util.replace_connection_catalog(
870            connection_id=self.connection_id,
871            configured_catalog_dict=catalog,
872            api_root=self.workspace.api_root,
873            client_id=self.workspace.client_id,
874            client_secret=self.workspace.client_secret,
875            bearer_token=self.workspace.bearer_token,
876            config_api_root=self.workspace.config_api_root,
877        )

Replace the configured catalog for this connection.

⚠️ WARNING: Modifying the catalog directly is not recommended and could result in broken connections, and/or incorrect sync behavior.

Accepts a configured catalog dict and replaces the connection's entire catalog with it. All other connection settings remain unchanged.

Accepts either format:

  • Config API format (syncCatalog with camelCase keys and nested config): passed through directly.
  • Airbyte protocol format (ConfiguredAirbyteCatalog with snake_case keys): automatically converted to Config API format before sending.
Arguments:
  • catalog: The configured catalog dict in either format.
def rename(self, name: str) -> CloudConnection:
879    def rename(self, name: str) -> CloudConnection:
880        """Rename the connection.
881
882        Args:
883            name: New name for the connection
884
885        Returns:
886            Updated CloudConnection object with refreshed info
887        """
888        updated_response = api_util.patch_connection(
889            connection_id=self.connection_id,
890            api_root=self.workspace.api_root,
891            client_id=self.workspace.client_id,
892            client_secret=self.workspace.client_secret,
893            bearer_token=self.workspace.bearer_token,
894            name=name,
895        )
896        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
897        return self

Rename the connection.

Arguments:
  • name: New name for the connection
Returns:

Updated CloudConnection object with refreshed info

def set_table_prefix(self, prefix: str) -> CloudConnection:
899    def set_table_prefix(self, prefix: str) -> CloudConnection:
900        """Set the table prefix for the connection.
901
902        Args:
903            prefix: New table prefix to use when syncing to the destination
904
905        Returns:
906            Updated CloudConnection object with refreshed info
907        """
908        updated_response = api_util.patch_connection(
909            connection_id=self.connection_id,
910            api_root=self.workspace.api_root,
911            client_id=self.workspace.client_id,
912            client_secret=self.workspace.client_secret,
913            bearer_token=self.workspace.bearer_token,
914            prefix=prefix,
915        )
916        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
917        return self

Set the table prefix for the connection.

Arguments:
  • prefix: New table prefix to use when syncing to the destination
Returns:

Updated CloudConnection object with refreshed info

def set_selected_streams( self, stream_names: list[str]) -> CloudConnection:
919    def set_selected_streams(self, stream_names: list[str]) -> CloudConnection:
920        """Set the selected streams for the connection.
921
922        This is a destructive operation that can break existing connections if the
923        stream selection is changed incorrectly. Use with caution.
924
925        Args:
926            stream_names: List of stream names to sync
927
928        Returns:
929            Updated CloudConnection object with refreshed info
930        """
931        configurations = api_util.build_stream_configurations(stream_names)
932
933        updated_response = api_util.patch_connection(
934            connection_id=self.connection_id,
935            api_root=self.workspace.api_root,
936            client_id=self.workspace.client_id,
937            client_secret=self.workspace.client_secret,
938            bearer_token=self.workspace.bearer_token,
939            configurations=configurations,
940        )
941        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)
942        return self

Set the selected streams for the connection.

This is a destructive operation that can break existing connections if the stream selection is changed incorrectly. Use with caution.

Arguments:
  • stream_names: List of stream names to sync
Returns:

Updated CloudConnection object with refreshed info

enabled: bool
946    @property
947    def enabled(self) -> bool:
948        """Get the current enabled status of the connection.
949
950        This property always fetches fresh data from the API to ensure accuracy,
951        as another process or user may have toggled the setting.
952
953        Returns:
954            True if the connection status is 'active', False otherwise.
955        """
956        connection_info = self._fetch_connection_info(force_refresh=True)
957        return connection_info.status == "active"

Get the current enabled status of the connection.

This property always fetches fresh data from the API to ensure accuracy, as another process or user may have toggled the setting.

Returns:

True if the connection status is 'active', False otherwise.

def set_enabled(self, *, enabled: bool, ignore_noop: bool = True) -> None:
 969    def set_enabled(
 970        self,
 971        *,
 972        enabled: bool,
 973        ignore_noop: bool = True,
 974    ) -> None:
 975        """Set the enabled status of the connection.
 976
 977        Args:
 978            enabled: True to enable (set status to 'active'), False to disable
 979                (set status to 'inactive').
 980            ignore_noop: If True (default), silently return if the connection is already
 981                in the requested state. If False, raise ValueError when the requested
 982                state matches the current state.
 983
 984        Raises:
 985            ValueError: If ignore_noop is False and the connection is already in the
 986                requested state.
 987        """
 988        # Always fetch fresh data to check current status
 989        connection_info = self._fetch_connection_info(force_refresh=True)
 990        current_status = connection_info.status
 991        desired_status = "active" if enabled else "inactive"
 992
 993        if current_status == desired_status:
 994            if ignore_noop:
 995                return
 996            raise ValueError(
 997                f"Connection is already {'enabled' if enabled else 'disabled'}. "
 998                f"Current status: {current_status}"
 999            )
1000
1001        updated_response = api_util.patch_connection(
1002            connection_id=self.connection_id,
1003            api_root=self.workspace.api_root,
1004            client_id=self.workspace.client_id,
1005            client_secret=self.workspace.client_secret,
1006            bearer_token=self.workspace.bearer_token,
1007            status=desired_status,
1008        )
1009        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)

Set the enabled status of the connection.

Arguments:
  • enabled: True to enable (set status to 'active'), False to disable (set status to 'inactive').
  • ignore_noop: If True (default), silently return if the connection is already in the requested state. If False, raise ValueError when the requested state matches the current state.
Raises:
  • ValueError: If ignore_noop is False and the connection is already in the requested state.
def set_schedule(self, cron_expression: str) -> None:
1013    def set_schedule(
1014        self,
1015        cron_expression: str,
1016    ) -> None:
1017        """Set a cron schedule for the connection.
1018
1019        Args:
1020            cron_expression: A Quartz cron expression defining when syncs should run.
1021                Quartz expressions have 6 or 7 space-separated fields
1022                (seconds, minutes, hours, day-of-month, month, day-of-week[, year]),
1023                optionally followed by a timezone ID. The Airbyte API rejects standard
1024                5-field Unix cron expressions and schedules that run more often than
1025                once per hour.
1026
1027        Examples:
1028            - "0 0 0 * * ?"  # Daily at midnight UTC
1029            - "0 0 */6 * * ?"  # Every 6 hours
1030            - "0 0 0 ? * SUN"  # Weekly on Sunday at midnight UTC
1031            - "0 0 9 ? * MON-FRI US/Pacific"  # Weekdays at 9am Pacific
1032        """
1033        _validate_quartz_cron_expression(cron_expression)
1034        updated_response = api_util.patch_connection(
1035            connection_id=self.connection_id,
1036            api_root=self.workspace.api_root,
1037            client_id=self.workspace.client_id,
1038            client_secret=self.workspace.client_secret,
1039            bearer_token=self.workspace.bearer_token,
1040            schedule=api_util.build_connection_schedule(
1041                schedule_type="cron",
1042                cron_expression=cron_expression,
1043            ),
1044        )
1045        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)

Set a cron schedule for the connection.

Arguments:
  • cron_expression: A Quartz cron expression defining when syncs should run. Quartz expressions have 6 or 7 space-separated fields (seconds, minutes, hours, day-of-month, month, day-of-week[, year]), optionally followed by a timezone ID. The Airbyte API rejects standard 5-field Unix cron expressions and schedules that run more often than once per hour.
Examples:
  • "0 0 0 * * ?" # Daily at midnight UTC
  • "0 0 */6 * * ?" # Every 6 hours
  • "0 0 0 ? * SUN" # Weekly on Sunday at midnight UTC
  • "0 0 9 ? * MON-FRI US/Pacific" # Weekdays at 9am Pacific
def set_manual_schedule(self) -> None:
1047    def set_manual_schedule(self) -> None:
1048        """Set the connection to manual scheduling.
1049
1050        Disables automatic syncs. Syncs will only run when manually triggered.
1051        """
1052        updated_response = api_util.patch_connection(
1053            connection_id=self.connection_id,
1054            api_root=self.workspace.api_root,
1055            client_id=self.workspace.client_id,
1056            client_secret=self.workspace.client_secret,
1057            bearer_token=self.workspace.bearer_token,
1058            schedule=api_util.build_connection_schedule(schedule_type="manual"),
1059        )
1060        self._connection_info = CloudConnectionInfo.from_api_response(updated_response)

Set the connection to manual scheduling.

Disables automatic syncs. Syncs will only run when manually triggered.

def permanently_delete( self, *, cascade_delete_source: bool = False, cascade_delete_destination: bool = False) -> None:
1064    def permanently_delete(
1065        self,
1066        *,
1067        cascade_delete_source: bool = False,
1068        cascade_delete_destination: bool = False,
1069    ) -> None:
1070        """Delete the connection.
1071
1072        Args:
1073            cascade_delete_source: Whether to also delete the source.
1074            cascade_delete_destination: Whether to also delete the destination.
1075        """
1076        self.workspace.permanently_delete_connection(self)
1077
1078        if cascade_delete_source:
1079            self.workspace.permanently_delete_source(self.source_id)
1080
1081        if cascade_delete_destination:
1082            self.workspace.permanently_delete_destination(self.destination_id)

Delete the connection.

Arguments:
  • cascade_delete_source: Whether to also delete the source.
  • cascade_delete_destination: Whether to also delete the destination.
@dataclass
class CloudClientConfig:
 59@dataclass
 60class CloudClientConfig:
 61    """Client configuration for Airbyte Cloud API.
 62
 63    This class encapsulates the authentication and API configuration needed to connect
 64    to Airbyte Cloud, OSS, or Enterprise instances. It supports two mutually
 65    exclusive authentication methods:
 66
 67    1. OAuth2 client credentials flow (client_id + client_secret)
 68    2. Bearer token authentication
 69
 70    Exactly one authentication method must be provided. Providing both or neither
 71    will raise a validation error.
 72
 73    Attributes:
 74        client_id: OAuth2 client ID for client credentials flow.
 75        client_secret: OAuth2 client secret for client credentials flow.
 76        bearer_token: Pre-generated bearer token for direct authentication.
 77        api_root: The API root URL. Defaults to Airbyte Cloud API.
 78        config_api_root: The Config API root URL.
 79    """
 80
 81    client_id: SecretString | None = None
 82    """OAuth2 client ID for client credentials authentication."""
 83
 84    client_secret: SecretString | None = None
 85    """OAuth2 client secret for client credentials authentication."""
 86
 87    bearer_token: SecretString | None = None
 88    """Bearer token for direct authentication (alternative to client credentials)."""
 89
 90    api_root: str = api_util.CLOUD_API_ROOT
 91    """The API root URL. Defaults to Airbyte Cloud API."""
 92
 93    config_api_root: str | None = None
 94    """The Config API root URL."""
 95
 96    def __post_init__(self) -> None:
 97        """Validate credentials and ensure secrets are properly wrapped."""
 98        # Wrap secrets in SecretString if they aren't already
 99        if self.client_id is not None:
100            self.client_id = SecretString(self.client_id)
101        if self.client_secret is not None:
102            self.client_secret = SecretString(self.client_secret)
103        if self.bearer_token is not None:
104            self.bearer_token = SecretString(self.bearer_token)
105
106        # Validate mutual exclusivity
107        has_client_credentials = self.client_id is not None or self.client_secret is not None
108        has_bearer_token = self.bearer_token is not None
109
110        if has_client_credentials and has_bearer_token:
111            raise PyAirbyteInputError(
112                message="Cannot use both client credentials and bearer token authentication.",
113                guidance=(
114                    "Provide either client_id and client_secret together, "
115                    "or bearer_token alone, but not both."
116                ),
117            )
118
119        if has_client_credentials and (self.client_id is None or self.client_secret is None):
120            # If using client credentials, both must be provided
121            raise PyAirbyteInputError(
122                message="Incomplete client credentials.",
123                guidance=(
124                    "When using client credentials authentication, "
125                    "both client_id and client_secret must be provided."
126                ),
127            )
128
129        if not has_client_credentials and not has_bearer_token:
130            raise PyAirbyteInputError(
131                message="No authentication credentials provided.",
132                guidance=(
133                    "Provide either client_id and client_secret together for OAuth2 "
134                    "client credentials flow, or bearer_token for direct authentication."
135                ),
136            )
137
138    @property
139    def uses_bearer_token(self) -> bool:
140        """Return True if using bearer token authentication."""
141        return self.bearer_token is not None
142
143    @property
144    def uses_client_credentials(self) -> bool:
145        """Return True if using client credentials authentication."""
146        return self.client_id is not None and self.client_secret is not None
147
148    @classmethod
149    def from_env(
150        cls,
151        *,
152        api_root: str | None = None,
153        config_api_root: str | None = None,
154    ) -> CloudClientConfig:
155        """Create CloudClientConfig from environment variables.
156
157        This factory method resolves credentials from environment variables,
158        providing a convenient way to create credentials without explicitly
159        passing secrets.
160
161        Environment variables used:
162            - `AIRBYTE_CLOUD_CLIENT_ID`: OAuth client ID (for client credentials flow).
163            - `AIRBYTE_CLOUD_CLIENT_SECRET`: OAuth client secret (for client credentials flow).
164            - `AIRBYTE_CLOUD_BEARER_TOKEN`: Bearer token (alternative to client credentials).
165            - `AIRBYTE_CLOUD_API_URL`: Optional. The API root URL (defaults to Airbyte Cloud).
166            - `AIRBYTE_CLOUD_CONFIG_API_URL`: Optional. The Config API root URL.
167
168        The method will first check for a bearer token. If not found, it will
169        attempt to use client credentials.
170
171        Args:
172            api_root: The API root URL. If not provided, will be resolved from
173                the `AIRBYTE_CLOUD_API_URL` environment variable, or default to
174                the Airbyte Cloud API.
175            config_api_root: The Config API root URL. If not provided, will be resolved
176                from the `AIRBYTE_CLOUD_CONFIG_API_URL` environment variable.
177
178        Returns:
179            A CloudClientConfig instance configured with credentials from the environment.
180
181        Raises:
182            PyAirbyteSecretNotFoundError: If required credentials are not found in
183                the environment.
184        """
185        resolved_api_root = resolve_cloud_api_url(api_root)
186        resolved_config_api_root = resolve_cloud_config_api_url(config_api_root)
187
188        # Try bearer token first
189        bearer_token = resolve_cloud_bearer_token()
190        if bearer_token:
191            return cls(
192                bearer_token=bearer_token,
193                api_root=resolved_api_root,
194                config_api_root=resolved_config_api_root,
195            )
196
197        # Fall back to client credentials
198        return cls(
199            client_id=resolve_cloud_client_id(),
200            client_secret=resolve_cloud_client_secret(),
201            api_root=resolved_api_root,
202            config_api_root=resolved_config_api_root,
203        )

Client configuration for Airbyte Cloud API.

This class encapsulates the authentication and API configuration needed to connect to Airbyte Cloud, OSS, or Enterprise instances. It supports two mutually exclusive authentication methods:

  1. OAuth2 client credentials flow (client_id + client_secret)
  2. Bearer token authentication

Exactly one authentication method must be provided. Providing both or neither will raise a validation error.

Attributes:
  • client_id: OAuth2 client ID for client credentials flow.
  • client_secret: OAuth2 client secret for client credentials flow.
  • bearer_token: Pre-generated bearer token for direct authentication.
  • api_root: The API root URL. Defaults to Airbyte Cloud API.
  • config_api_root: The Config API root URL.
CloudClientConfig( client_id: airbyte.secrets.SecretString | None = None, client_secret: airbyte.secrets.SecretString | None = None, bearer_token: airbyte.secrets.SecretString | None = None, api_root: str = 'https://api.airbyte.com/v1', config_api_root: str | None = None)
client_id: airbyte.secrets.SecretString | None = None

OAuth2 client ID for client credentials authentication.

client_secret: airbyte.secrets.SecretString | None = None

OAuth2 client secret for client credentials authentication.

bearer_token: airbyte.secrets.SecretString | None = None

Bearer token for direct authentication (alternative to client credentials).

api_root: str = 'https://api.airbyte.com/v1'

The API root URL. Defaults to Airbyte Cloud API.

config_api_root: str | None = None

The Config API root URL.

uses_bearer_token: bool
138    @property
139    def uses_bearer_token(self) -> bool:
140        """Return True if using bearer token authentication."""
141        return self.bearer_token is not None

Return True if using bearer token authentication.

uses_client_credentials: bool
143    @property
144    def uses_client_credentials(self) -> bool:
145        """Return True if using client credentials authentication."""
146        return self.client_id is not None and self.client_secret is not None

Return True if using client credentials authentication.

@classmethod
def from_env( cls, *, api_root: str | None = None, config_api_root: str | None = None) -> CloudClientConfig:
148    @classmethod
149    def from_env(
150        cls,
151        *,
152        api_root: str | None = None,
153        config_api_root: str | None = None,
154    ) -> CloudClientConfig:
155        """Create CloudClientConfig from environment variables.
156
157        This factory method resolves credentials from environment variables,
158        providing a convenient way to create credentials without explicitly
159        passing secrets.
160
161        Environment variables used:
162            - `AIRBYTE_CLOUD_CLIENT_ID`: OAuth client ID (for client credentials flow).
163            - `AIRBYTE_CLOUD_CLIENT_SECRET`: OAuth client secret (for client credentials flow).
164            - `AIRBYTE_CLOUD_BEARER_TOKEN`: Bearer token (alternative to client credentials).
165            - `AIRBYTE_CLOUD_API_URL`: Optional. The API root URL (defaults to Airbyte Cloud).
166            - `AIRBYTE_CLOUD_CONFIG_API_URL`: Optional. The Config API root URL.
167
168        The method will first check for a bearer token. If not found, it will
169        attempt to use client credentials.
170
171        Args:
172            api_root: The API root URL. If not provided, will be resolved from
173                the `AIRBYTE_CLOUD_API_URL` environment variable, or default to
174                the Airbyte Cloud API.
175            config_api_root: The Config API root URL. If not provided, will be resolved
176                from the `AIRBYTE_CLOUD_CONFIG_API_URL` environment variable.
177
178        Returns:
179            A CloudClientConfig instance configured with credentials from the environment.
180
181        Raises:
182            PyAirbyteSecretNotFoundError: If required credentials are not found in
183                the environment.
184        """
185        resolved_api_root = resolve_cloud_api_url(api_root)
186        resolved_config_api_root = resolve_cloud_config_api_url(config_api_root)
187
188        # Try bearer token first
189        bearer_token = resolve_cloud_bearer_token()
190        if bearer_token:
191            return cls(
192                bearer_token=bearer_token,
193                api_root=resolved_api_root,
194                config_api_root=resolved_config_api_root,
195            )
196
197        # Fall back to client credentials
198        return cls(
199            client_id=resolve_cloud_client_id(),
200            client_secret=resolve_cloud_client_secret(),
201            api_root=resolved_api_root,
202            config_api_root=resolved_config_api_root,
203        )

Create CloudClientConfig from environment variables.

This factory method resolves credentials from environment variables, providing a convenient way to create credentials without explicitly passing secrets.

Environment variables used:
  • AIRBYTE_CLOUD_CLIENT_ID: OAuth client ID (for client credentials flow).
  • AIRBYTE_CLOUD_CLIENT_SECRET: OAuth client secret (for client credentials flow).
  • AIRBYTE_CLOUD_BEARER_TOKEN: Bearer token (alternative to client credentials).
  • AIRBYTE_CLOUD_API_URL: Optional. The API root URL (defaults to Airbyte Cloud).
  • AIRBYTE_CLOUD_CONFIG_API_URL: Optional. The Config API root URL.

The method will first check for a bearer token. If not found, it will attempt to use client credentials.

Arguments:
  • api_root: The API root URL. If not provided, will be resolved from the AIRBYTE_CLOUD_API_URL environment variable, or default to the Airbyte Cloud API.
  • config_api_root: The Config API root URL. If not provided, will be resolved from the AIRBYTE_CLOUD_CONFIG_API_URL environment variable.
Returns:

A CloudClientConfig instance configured with credentials from the environment.

Raises:
  • PyAirbyteSecretNotFoundError: If required credentials are not found in the environment.
class CloudDefaultContextInfo(pydantic.main.BaseModel):
192class CloudDefaultContextInfo(BaseModel):
193    """Explicit organization and workspace affinities for the authenticated user."""
194
195    user_id: str | None
196    """The Airbyte user ID, if available."""
197
198    user_name: str | None
199    """The authenticated user's name, if available."""
200
201    user_email: str | None
202    """The authenticated user's email, if available."""
203
204    default_workspace_id: str | None
205    """The resolved default workspace ID, if available."""
206
207    default_workspace_name: str | None
208    """The resolved default workspace name, if available."""
209
210    default_workspace_verified: bool
211    """Whether the resolved default workspace was verified as accessible."""
212
213    unvalidated_workspace_count: int = 0
214    """Number of direct workspace grants not validated due to the validation cap."""
215
216    default_organization_id: str | None
217    """The organization containing the resolved default workspace, if available."""
218
219    default_organization_name: str | None
220    """The name of the organization containing the resolved default workspace, if available."""
221
222    configured_workspace_id: str | None
223    """The explicitly configured workspace ID, if available."""
224
225    configured_organization_id: str | None
226    """The configured organization ID, if available."""
227
228    member_organizations: list[CloudOrganizationInfo]
229    """Organizations identified by explicit organization membership grants."""
230
231    member_workspaces: list[CloudWorkspaceInfo]
232    """Workspaces identified by explicit workspace membership grants."""
233
234    member_organizations_truncated: bool
235    """True if organization memberships beyond the returned list were omitted."""
236
237    member_workspaces_truncated: bool
238    """True if workspace memberships beyond the returned list were omitted."""
239
240    discovery_hints: list[str]
241    """Hints for discovering additional organizations or workspaces."""

Explicit organization and workspace affinities for the authenticated user.

user_id: str | None = PydanticUndefined

The Airbyte user ID, if available.

user_name: str | None = PydanticUndefined

The authenticated user's name, if available.

user_email: str | None = PydanticUndefined

The authenticated user's email, if available.

default_workspace_id: str | None = PydanticUndefined

The resolved default workspace ID, if available.

default_workspace_name: str | None = PydanticUndefined

The resolved default workspace name, if available.

default_workspace_verified: bool = PydanticUndefined

Whether the resolved default workspace was verified as accessible.

unvalidated_workspace_count: int = 0

Number of direct workspace grants not validated due to the validation cap.

default_organization_id: str | None = PydanticUndefined

The organization containing the resolved default workspace, if available.

default_organization_name: str | None = PydanticUndefined

The name of the organization containing the resolved default workspace, if available.

configured_workspace_id: str | None = PydanticUndefined

The explicitly configured workspace ID, if available.

configured_organization_id: str | None = PydanticUndefined

The configured organization ID, if available.

member_organizations: list[airbyte.cloud.models.CloudOrganizationInfo] = PydanticUndefined

Organizations identified by explicit organization membership grants.

member_workspaces: list[CloudWorkspaceInfo] = PydanticUndefined

Workspaces identified by explicit workspace membership grants.

member_organizations_truncated: bool = PydanticUndefined

True if organization memberships beyond the returned list were omitted.

member_workspaces_truncated: bool = PydanticUndefined

True if workspace memberships beyond the returned list were omitted.

discovery_hints: list[str] = PydanticUndefined

Hints for discovering additional organizations or workspaces.

class CloudWorkspaceInfo(pydantic.main.BaseModel):
134class CloudWorkspaceInfo(BaseModel):
135    """Information about an Airbyte workspace."""
136
137    model_config = ConfigDict(populate_by_name=True)
138
139    workspace_id: str = Field(alias="workspaceId")
140    """The workspace ID."""
141
142    name: str
143    """The workspace name."""
144
145    data_residency: str | None = Field(default=None, alias="dataResidency")
146    """The data residency setting for the workspace, if available."""
147
148    organization_id: str | None = Field(default=None, alias="organizationId")
149    """The organization ID for the workspace, if available."""
150
151    organization_name: str | None = Field(default=None, alias="organizationName")
152    """The organization name for the workspace, if available."""
153
154    notifications: dict[str, object | None] | list[dict[str, object | None]] = Field(
155        default_factory=dict
156    )
157    """Workspace notification settings."""
158
159    @classmethod
160    def from_api_response(cls, workspace: _WorkspaceResponseLike) -> CloudWorkspaceInfo:
161        """Create a public model from an internal API workspace response."""
162        return cls(
163            workspace_id=workspace.workspace_id,
164            name=workspace.name,
165            data_residency=workspace.data_residency,
166            organization_id=getattr(workspace, "organization_id", None),
167            notifications=_notifications_to_dict(workspace.notifications),
168        )
169
170    @classmethod
171    def from_mapping(cls, workspace: Mapping[str, object]) -> CloudWorkspaceInfo:
172        """Create a public model from a workspace mapping."""
173        return cls.model_validate(workspace)
174
175    def to_dict(self) -> dict[str, object]:
176        """Return a JSON-serializable dictionary."""
177        return self.model_dump(mode="json")

Information about an Airbyte workspace.

workspace_id: str = PydanticUndefined

The workspace ID.

name: str = PydanticUndefined

The workspace name.

data_residency: str | None = None

The data residency setting for the workspace, if available.

organization_id: str | None = None

The organization ID for the workspace, if available.

organization_name: str | None = None

The organization name for the workspace, if available.

notifications: dict[str, object | None] | list[dict[str, object | None]] = PydanticUndefined

Workspace notification settings.

@classmethod
def from_api_response( cls, workspace: airbyte.cloud.models._WorkspaceResponseLike) -> CloudWorkspaceInfo:
159    @classmethod
160    def from_api_response(cls, workspace: _WorkspaceResponseLike) -> CloudWorkspaceInfo:
161        """Create a public model from an internal API workspace response."""
162        return cls(
163            workspace_id=workspace.workspace_id,
164            name=workspace.name,
165            data_residency=workspace.data_residency,
166            organization_id=getattr(workspace, "organization_id", None),
167            notifications=_notifications_to_dict(workspace.notifications),
168        )

Create a public model from an internal API workspace response.

@classmethod
def from_mapping( cls, workspace: Mapping[str, object]) -> CloudWorkspaceInfo:
170    @classmethod
171    def from_mapping(cls, workspace: Mapping[str, object]) -> CloudWorkspaceInfo:
172        """Create a public model from a workspace mapping."""
173        return cls.model_validate(workspace)

Create a public model from a workspace mapping.

def to_dict(self) -> dict[str, object]:
175    def to_dict(self) -> dict[str, object]:
176        """Return a JSON-serializable dictionary."""
177        return self.model_dump(mode="json")

Return a JSON-serializable dictionary.

@dataclass
class SyncResult:
218@dataclass
219class SyncResult:
220    """The result of a sync operation.
221
222    **This class is not meant to be instantiated directly.** Instead, obtain a `SyncResult` by
223    interacting with the `.CloudWorkspace` and `.CloudConnection` objects.
224    """
225
226    workspace: CloudWorkspace
227    connection: CloudConnection
228    job_id: int
229    table_name_prefix: str = ""
230    table_name_suffix: str = ""
231    _latest_job_info: CloudJobInfo | None = None
232    _connection_response: CloudConnectionInfo | None = None
233    _cache: CacheBase | None = None
234    _job_with_attempts_info: dict[str, Any] | None = None
235
236    @property
237    def job_url(self) -> str:
238        """Return the URL of the sync job.
239
240        Note: This currently returns the connection's job history URL, as there is no direct URL
241        to a specific job in the Airbyte Cloud web app.
242
243        TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number.
244              E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true
245        """
246        return f"{self.connection.job_history_url}"
247
248    def _get_connection_info(self, *, force_refresh: bool = False) -> CloudConnectionInfo:
249        """Return connection info for the sync job."""
250        if self._connection_response and not force_refresh:
251            return self._connection_response
252
253        self._connection_response = CloudConnectionInfo.from_api_response(
254            api_util.get_connection(
255                workspace_id=self.workspace.workspace_id,
256                api_root=self.workspace.api_root,
257                connection_id=self.connection.connection_id,
258                client_id=self.workspace.client_id,
259                client_secret=self.workspace.client_secret,
260                bearer_token=self.workspace.bearer_token,
261            )
262        )
263        return self._connection_response
264
265    def _get_destination_configuration(self, *, force_refresh: bool = False) -> dict[str, Any]:
266        """Return the destination configuration for the sync job."""
267        connection_info = self._get_connection_info(force_refresh=force_refresh)
268        destination_response = api_util.get_destination(
269            destination_id=connection_info.destination_id,
270            api_root=self.workspace.api_root,
271            client_id=self.workspace.client_id,
272            client_secret=self.workspace.client_secret,
273            bearer_token=self.workspace.bearer_token,
274        )
275        configuration = destination_response.configuration
276        if isinstance(configuration, Mapping):
277            configuration_dict = configuration
278        else:
279            configuration_dict = asdict(configuration)
280        return {**configuration_dict, "destinationType": destination_response.destination_type}
281
282    def is_job_complete(self) -> bool:
283        """Check if the sync job is complete."""
284        return self.get_job_status() in FINAL_STATUSES
285
286    def get_job_status(self) -> JobStatusEnum:
287        """Check if the sync job is still running."""
288        return self._fetch_latest_job_info().status
289
290    def _fetch_latest_job_info(self) -> CloudJobInfo:
291        """Return the job info for the sync job."""
292        if self._latest_job_info and self._latest_job_info.status in FINAL_STATUSES:
293            return self._latest_job_info
294
295        self._latest_job_info = CloudJobInfo.from_api_response(
296            api_util.get_job_info(
297                job_id=self.job_id,
298                api_root=self.workspace.api_root,
299                client_id=self.workspace.client_id,
300                client_secret=self.workspace.client_secret,
301                bearer_token=self.workspace.bearer_token,
302            )
303        )
304        return self._latest_job_info
305
306    @property
307    def bytes_synced(self) -> int:
308        """Return the number of records processed."""
309        return self._fetch_latest_job_info().bytes_synced or 0
310
311    @property
312    def records_synced(self) -> int:
313        """Return the number of records processed."""
314        return self._fetch_latest_job_info().rows_synced or 0
315
316    @property
317    def start_time(self) -> datetime:
318        """Return the start time of the sync job in UTC."""
319        try:
320            return ab_datetime_parse(self._fetch_latest_job_info().start_time)
321        except (ValueError, TypeError) as e:
322            if "Invalid isoformat string" in str(e):
323                job_info_raw = api_util._make_config_api_request(  # noqa: SLF001
324                    api_root=self.workspace.api_root,
325                    config_api_root=self.workspace.config_api_root,
326                    path="/jobs/get",
327                    json={"id": self.job_id},
328                    client_id=self.workspace.client_id,
329                    client_secret=self.workspace.client_secret,
330                    bearer_token=self.workspace.bearer_token,
331                )
332                raw_start_time = job_info_raw.get("startTime")
333                if raw_start_time:
334                    return ab_datetime_parse(raw_start_time)
335            raise
336
337    def _fetch_job_with_attempts(self) -> dict[str, Any]:
338        """Fetch job info with attempts from Config API using lazy loading pattern."""
339        if self._job_with_attempts_info is not None:
340            return self._job_with_attempts_info
341
342        self._job_with_attempts_info = api_util._make_config_api_request(  # noqa: SLF001  # Config API helper
343            api_root=self.workspace.api_root,
344            config_api_root=self.workspace.config_api_root,
345            path="/jobs/get",
346            json={
347                "id": self.job_id,
348            },
349            client_id=self.workspace.client_id,
350            client_secret=self.workspace.client_secret,
351            bearer_token=self.workspace.bearer_token,
352        )
353        return self._job_with_attempts_info
354
355    def get_attempts(self) -> list[SyncAttempt]:
356        """Return a list of attempts for this sync job."""
357        job_with_attempts = self._fetch_job_with_attempts()
358        attempts_data = job_with_attempts.get("attempts", [])
359
360        return [
361            SyncAttempt(
362                workspace=self.workspace,
363                connection=self.connection,
364                job_id=self.job_id,
365                attempt_number=i,
366                _attempt_data=attempt_data,
367            )
368            for i, attempt_data in enumerate(attempts_data, start=0)
369        ]
370
371    def raise_failure_status(
372        self,
373        *,
374        refresh_status: bool = False,
375    ) -> None:
376        """Raise an exception if the sync job failed.
377
378        By default, this method will use the latest status available. If you want to refresh the
379        status before checking for failure, set `refresh_status=True`. If the job has failed, this
380        method will raise a `AirbyteConnectionSyncError`.
381
382        Otherwise, do nothing.
383        """
384        if not refresh_status and self._latest_job_info:
385            latest_status = self._latest_job_info.status
386        else:
387            latest_status = self.get_job_status()
388
389        if latest_status in FAILED_STATUSES:
390            raise AirbyteConnectionSyncError(
391                workspace=self.workspace,
392                connection_id=self.connection.connection_id,
393                job_id=self.job_id,
394                job_status=self.get_job_status(),
395            )
396
397    def wait_for_completion(
398        self,
399        *,
400        wait_timeout: int = DEFAULT_SYNC_TIMEOUT_SECONDS,
401        raise_timeout: bool = True,
402        raise_failure: bool = False,
403    ) -> JobStatusEnum:
404        """Wait for a job to finish running."""
405        start_time = time.time()
406        while True:
407            latest_status = self.get_job_status()
408            if latest_status in FINAL_STATUSES:
409                if raise_failure:
410                    # No-op if the job succeeded or is still running:
411                    self.raise_failure_status()
412
413                return latest_status
414
415            if time.time() - start_time > wait_timeout:
416                if raise_timeout:
417                    raise AirbyteConnectionSyncTimeoutError(
418                        workspace=self.workspace,
419                        connection_id=self.connection.connection_id,
420                        job_id=self.job_id,
421                        job_status=latest_status,
422                        timeout=wait_timeout,
423                    )
424
425                return latest_status  # This will be a non-final status
426
427            time.sleep(api_util.JOB_WAIT_INTERVAL_SECS)
428
429    def get_sql_cache(self) -> CacheBase:
430        """Return a SQL Cache object for working with the data in a SQL-based destination's."""
431        if self._cache:
432            return self._cache
433
434        destination_configuration = self._get_destination_configuration()
435        self._cache = destination_to_cache(destination_configuration=destination_configuration)
436        return self._cache
437
438    def get_sql_engine(self) -> sqlalchemy.engine.Engine:
439        """Return a SQL Engine for querying a SQL-based destination."""
440        return self.get_sql_cache().get_sql_engine()
441
442    def get_sql_table_name(self, stream_name: str) -> str:
443        """Return the SQL table name of the named stream."""
444        return self.get_sql_cache().processor.get_sql_table_name(stream_name=stream_name)
445
446    def get_sql_table(
447        self,
448        stream_name: str,
449    ) -> sqlalchemy.Table:
450        """Return a SQLAlchemy table object for the named stream."""
451        return self.get_sql_cache().processor.get_sql_table(stream_name)
452
453    def get_dataset(self, stream_name: str) -> CachedDataset:
454        """Retrieve an `airbyte.datasets.CachedDataset` object for a given stream name.
455
456        This can be used to read and analyze the data in a SQL-based destination.
457
458        TODO: In a future iteration, we can consider providing stream configuration information
459              (catalog information) to the `CachedDataset` object via the "Get stream properties"
460              API: https://reference.airbyte.com/reference/getstreamproperties
461        """
462        return CachedDataset(
463            self.get_sql_cache(),
464            stream_name=stream_name,
465            stream_configuration=False,  # Don't look for stream configuration in cache.
466        )
467
468    def get_sql_database_name(self) -> str:
469        """Return the SQL database name."""
470        cache = self.get_sql_cache()
471        return cache.get_database_name()
472
473    def get_sql_schema_name(self) -> str:
474        """Return the SQL schema name."""
475        cache = self.get_sql_cache()
476        return cache.schema_name
477
478    @property
479    def stream_names(self) -> list[str]:
480        """Return the set of stream names."""
481        return self.connection.stream_names
482
483    @final
484    @property
485    def streams(
486        self,
487    ) -> _SyncResultStreams:  # pyrefly: ignore[unknown-name]
488        """Return a mapping of stream names to `airbyte.CachedDataset` objects.
489
490        This is a convenience wrapper around the `stream_names`
491        property and `get_dataset()` method.
492        """
493        return self._SyncResultStreams(self)
494
495    class _SyncResultStreams(Mapping[str, CachedDataset]):
496        """A mapping of stream names to cached datasets."""
497
498        def __init__(
499            self,
500            parent: SyncResult,
501            /,
502        ) -> None:
503            self.parent: SyncResult = parent
504
505        def __getitem__(self, key: str) -> CachedDataset:
506            return self.parent.get_dataset(stream_name=key)
507
508        def __iter__(self) -> Iterator[str]:
509            return iter(self.parent.stream_names)
510
511        def __len__(self) -> int:
512            return len(self.parent.stream_names)

The result of a sync operation.

This class is not meant to be instantiated directly. Instead, obtain a SyncResult by interacting with the .CloudWorkspace and .CloudConnection objects.

SyncResult( workspace: CloudWorkspace, connection: CloudConnection, job_id: int, table_name_prefix: str = '', table_name_suffix: str = '', _latest_job_info: airbyte.cloud.models.CloudJobInfo | None = None, _connection_response: airbyte.cloud.models.CloudConnectionInfo | None = None, _cache: airbyte.caches.CacheBase | None = None, _job_with_attempts_info: dict[str, typing.Any] | None = None)
workspace: CloudWorkspace
connection: CloudConnection
job_id: int
table_name_prefix: str = ''
table_name_suffix: str = ''
job_url: str
236    @property
237    def job_url(self) -> str:
238        """Return the URL of the sync job.
239
240        Note: This currently returns the connection's job history URL, as there is no direct URL
241        to a specific job in the Airbyte Cloud web app.
242
243        TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number.
244              E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true
245        """
246        return f"{self.connection.job_history_url}"

Return the URL of the sync job.

Note: This currently returns the connection's job history URL, as there is no direct URL to a specific job in the Airbyte Cloud web app.

TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number. E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true

def is_job_complete(self) -> bool:
282    def is_job_complete(self) -> bool:
283        """Check if the sync job is complete."""
284        return self.get_job_status() in FINAL_STATUSES

Check if the sync job is complete.

def get_job_status(self) -> JobStatusEnum:
286    def get_job_status(self) -> JobStatusEnum:
287        """Check if the sync job is still running."""
288        return self._fetch_latest_job_info().status

Check if the sync job is still running.

bytes_synced: int
306    @property
307    def bytes_synced(self) -> int:
308        """Return the number of records processed."""
309        return self._fetch_latest_job_info().bytes_synced or 0

Return the number of records processed.

records_synced: int
311    @property
312    def records_synced(self) -> int:
313        """Return the number of records processed."""
314        return self._fetch_latest_job_info().rows_synced or 0

Return the number of records processed.

start_time: datetime.datetime
316    @property
317    def start_time(self) -> datetime:
318        """Return the start time of the sync job in UTC."""
319        try:
320            return ab_datetime_parse(self._fetch_latest_job_info().start_time)
321        except (ValueError, TypeError) as e:
322            if "Invalid isoformat string" in str(e):
323                job_info_raw = api_util._make_config_api_request(  # noqa: SLF001
324                    api_root=self.workspace.api_root,
325                    config_api_root=self.workspace.config_api_root,
326                    path="/jobs/get",
327                    json={"id": self.job_id},
328                    client_id=self.workspace.client_id,
329                    client_secret=self.workspace.client_secret,
330                    bearer_token=self.workspace.bearer_token,
331                )
332                raw_start_time = job_info_raw.get("startTime")
333                if raw_start_time:
334                    return ab_datetime_parse(raw_start_time)
335            raise

Return the start time of the sync job in UTC.

def get_attempts(self) -> list[airbyte.cloud.sync_results.SyncAttempt]:
355    def get_attempts(self) -> list[SyncAttempt]:
356        """Return a list of attempts for this sync job."""
357        job_with_attempts = self._fetch_job_with_attempts()
358        attempts_data = job_with_attempts.get("attempts", [])
359
360        return [
361            SyncAttempt(
362                workspace=self.workspace,
363                connection=self.connection,
364                job_id=self.job_id,
365                attempt_number=i,
366                _attempt_data=attempt_data,
367            )
368            for i, attempt_data in enumerate(attempts_data, start=0)
369        ]

Return a list of attempts for this sync job.

def raise_failure_status(self, *, refresh_status: bool = False) -> None:
371    def raise_failure_status(
372        self,
373        *,
374        refresh_status: bool = False,
375    ) -> None:
376        """Raise an exception if the sync job failed.
377
378        By default, this method will use the latest status available. If you want to refresh the
379        status before checking for failure, set `refresh_status=True`. If the job has failed, this
380        method will raise a `AirbyteConnectionSyncError`.
381
382        Otherwise, do nothing.
383        """
384        if not refresh_status and self._latest_job_info:
385            latest_status = self._latest_job_info.status
386        else:
387            latest_status = self.get_job_status()
388
389        if latest_status in FAILED_STATUSES:
390            raise AirbyteConnectionSyncError(
391                workspace=self.workspace,
392                connection_id=self.connection.connection_id,
393                job_id=self.job_id,
394                job_status=self.get_job_status(),
395            )

Raise an exception if the sync job failed.

By default, this method will use the latest status available. If you want to refresh the status before checking for failure, set refresh_status=True. If the job has failed, this method will raise a AirbyteConnectionSyncError.

Otherwise, do nothing.

def wait_for_completion( self, *, wait_timeout: int = 1800, raise_timeout: bool = True, raise_failure: bool = False) -> JobStatusEnum:
397    def wait_for_completion(
398        self,
399        *,
400        wait_timeout: int = DEFAULT_SYNC_TIMEOUT_SECONDS,
401        raise_timeout: bool = True,
402        raise_failure: bool = False,
403    ) -> JobStatusEnum:
404        """Wait for a job to finish running."""
405        start_time = time.time()
406        while True:
407            latest_status = self.get_job_status()
408            if latest_status in FINAL_STATUSES:
409                if raise_failure:
410                    # No-op if the job succeeded or is still running:
411                    self.raise_failure_status()
412
413                return latest_status
414
415            if time.time() - start_time > wait_timeout:
416                if raise_timeout:
417                    raise AirbyteConnectionSyncTimeoutError(
418                        workspace=self.workspace,
419                        connection_id=self.connection.connection_id,
420                        job_id=self.job_id,
421                        job_status=latest_status,
422                        timeout=wait_timeout,
423                    )
424
425                return latest_status  # This will be a non-final status
426
427            time.sleep(api_util.JOB_WAIT_INTERVAL_SECS)

Wait for a job to finish running.

def get_sql_cache(self) -> airbyte.caches.CacheBase:
429    def get_sql_cache(self) -> CacheBase:
430        """Return a SQL Cache object for working with the data in a SQL-based destination's."""
431        if self._cache:
432            return self._cache
433
434        destination_configuration = self._get_destination_configuration()
435        self._cache = destination_to_cache(destination_configuration=destination_configuration)
436        return self._cache

Return a SQL Cache object for working with the data in a SQL-based destination's.

def get_sql_engine(self) -> sqlalchemy.engine.base.Engine:
438    def get_sql_engine(self) -> sqlalchemy.engine.Engine:
439        """Return a SQL Engine for querying a SQL-based destination."""
440        return self.get_sql_cache().get_sql_engine()

Return a SQL Engine for querying a SQL-based destination.

def get_sql_table_name(self, stream_name: str) -> str:
442    def get_sql_table_name(self, stream_name: str) -> str:
443        """Return the SQL table name of the named stream."""
444        return self.get_sql_cache().processor.get_sql_table_name(stream_name=stream_name)

Return the SQL table name of the named stream.

def get_sql_table(self, stream_name: str) -> sqlalchemy.sql.schema.Table:
446    def get_sql_table(
447        self,
448        stream_name: str,
449    ) -> sqlalchemy.Table:
450        """Return a SQLAlchemy table object for the named stream."""
451        return self.get_sql_cache().processor.get_sql_table(stream_name)

Return a SQLAlchemy table object for the named stream.

def get_dataset(self, stream_name: str) -> airbyte.CachedDataset:
453    def get_dataset(self, stream_name: str) -> CachedDataset:
454        """Retrieve an `airbyte.datasets.CachedDataset` object for a given stream name.
455
456        This can be used to read and analyze the data in a SQL-based destination.
457
458        TODO: In a future iteration, we can consider providing stream configuration information
459              (catalog information) to the `CachedDataset` object via the "Get stream properties"
460              API: https://reference.airbyte.com/reference/getstreamproperties
461        """
462        return CachedDataset(
463            self.get_sql_cache(),
464            stream_name=stream_name,
465            stream_configuration=False,  # Don't look for stream configuration in cache.
466        )

Retrieve an airbyte.datasets.CachedDataset object for a given stream name.

This can be used to read and analyze the data in a SQL-based destination.

TODO: In a future iteration, we can consider providing stream configuration information (catalog information) to the CachedDataset object via the "Get stream properties" API: https://reference.airbyte.com/reference/getstreamproperties

def get_sql_database_name(self) -> str:
468    def get_sql_database_name(self) -> str:
469        """Return the SQL database name."""
470        cache = self.get_sql_cache()
471        return cache.get_database_name()

Return the SQL database name.

def get_sql_schema_name(self) -> str:
473    def get_sql_schema_name(self) -> str:
474        """Return the SQL schema name."""
475        cache = self.get_sql_cache()
476        return cache.schema_name

Return the SQL schema name.

stream_names: list[str]
478    @property
479    def stream_names(self) -> list[str]:
480        """Return the set of stream names."""
481        return self.connection.stream_names

Return the set of stream names.

streams: airbyte.cloud.sync_results.SyncResult._SyncResultStreams
483    @final
484    @property
485    def streams(
486        self,
487    ) -> _SyncResultStreams:  # pyrefly: ignore[unknown-name]
488        """Return a mapping of stream names to `airbyte.CachedDataset` objects.
489
490        This is a convenience wrapper around the `stream_names`
491        property and `get_dataset()` method.
492        """
493        return self._SyncResultStreams(self)

Return a mapping of stream names to airbyte.CachedDataset objects.

This is a convenience wrapper around the stream_names property and get_dataset() method.

class JobStatusEnum(builtins.str, enum.Enum):
 92class JobStatusEnum(str, Enum):
 93    """Status values for an Airbyte Cloud job."""
 94
 95    PENDING = "pending"
 96    RUNNING = "running"
 97    INCOMPLETE = "incomplete"
 98    FAILED = "failed"
 99    SUCCEEDED = "succeeded"
100    CANCELLED = "cancelled"

Status values for an Airbyte Cloud job.

PENDING = <JobStatusEnum.PENDING: 'pending'>
RUNNING = <JobStatusEnum.RUNNING: 'running'>
INCOMPLETE = <JobStatusEnum.INCOMPLETE: 'incomplete'>
FAILED = <JobStatusEnum.FAILED: 'failed'>
SUCCEEDED = <JobStatusEnum.SUCCEEDED: 'succeeded'>
CANCELLED = <JobStatusEnum.CANCELLED: 'cancelled'>
class JobTypeEnum(builtins.str, enum.Enum):
103class JobTypeEnum(str, Enum):
104    """Job type values for Airbyte Cloud jobs."""
105
106    SYNC = "sync"
107    RESET = "reset"
108    REFRESH = "refresh"
109    CLEAR = "clear"

Job type values for Airbyte Cloud jobs.

SYNC = <JobTypeEnum.SYNC: 'sync'>
RESET = <JobTypeEnum.RESET: 'reset'>
REFRESH = <JobTypeEnum.REFRESH: 'refresh'>
CLEAR = <JobTypeEnum.CLEAR: 'clear'>
class WorkspacePrivilegeScope(builtins.str, enum.Enum):
112class WorkspacePrivilegeScope(str, Enum):
113    """How broadly `list_workspaces` searches for workspaces."""
114
115    MEMBER_OF = "member_of"
116    ORGANIZATION_ADMIN = "organization_admin"
117    INSTANCE_ADMIN = "instance_admin"
118    ANY = "any"

How broadly list_workspaces searches for workspaces.

MEMBER_OF = <WorkspacePrivilegeScope.MEMBER_OF: 'member_of'>
ORGANIZATION_ADMIN = <WorkspacePrivilegeScope.ORGANIZATION_ADMIN: 'organization_admin'>
INSTANCE_ADMIN = <WorkspacePrivilegeScope.INSTANCE_ADMIN: 'instance_admin'>