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