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

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

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

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

It is not recommended to create a CloudConnection object directly.

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

connection_id

The ID of the connection.

workspace

The workspace that the connection belongs to.

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

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

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

Returns:

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

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

Get the display name of the connection, if available.

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

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

The ID of the source.

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

Get the source object.

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

The ID of the destination.

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

Get the destination object.

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

The stream names.

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

The table prefix.

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

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

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

The namespace format template, when namespace_definition is custom_format.

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

The web URL to the connection.

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

The URL to the job history for the connection.

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

Run a sync.

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

Cancel a running sync job.

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

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

Get previous sync jobs for a connection with pagination support.

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

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

A list of SyncResult objects representing the sync jobs.

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

Get the sync result for the connection.

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

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

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

Deprecated. Use dump_raw_state() instead.

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

Dump the state for this connection.

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

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

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

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

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

Import (restore) the full state for this connection.

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

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

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

Accepts either format:

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

The updated connection state as a dictionary.

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

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

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

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

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

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

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

Set the state for a single stream within this connection.

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

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

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

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

Get the configured catalog for this connection.

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

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

Returns:

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

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

Dump the configured catalog for this connection.

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

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

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

The configured catalog dict, or None if not found.

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

Replace the configured catalog for this connection.

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

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

Accepts either format:

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

Rename the connection.

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

Updated CloudConnection object with refreshed info

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

Set the table prefix for the connection.

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

Updated CloudConnection object with refreshed info

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

Set the selected streams for the connection.

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

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

Updated CloudConnection object with refreshed info

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

Get the current enabled status of the connection.

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

Returns:

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

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

Set the enabled status of the connection.

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

Set a cron schedule for the connection.

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

Set the connection to manual scheduling.

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

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

Delete the connection.

Arguments:
  • cascade_delete_source: Whether to also delete the source.
  • cascade_delete_destination: Whether to also delete the destination.