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