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