airbyte.cloud.sync_results

Sync results for Airbyte Cloud workspaces.

Examples

Run a sync job and wait for completion

To get started, we'll need a .CloudConnection object. You can obtain this object by calling .CloudWorkspace.get_connection().

from airbyte import cloud

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

# Get a connection object
connection = workspace.get_connection(connection_id="456")

Once we have a .CloudConnection object, we can simply call run_sync() to start a sync job and wait for it to complete.

# Run a sync job
sync_result: SyncResult = connection.run_sync()

Run a sync job and return immediately

By default, run_sync() will wait for the job to complete and raise an exception if the job fails. You can instead return immediately by setting wait=False.

# Start the sync job and return immediately
sync_result: SyncResult = connection.run_sync(wait=False)

while not sync_result.is_job_complete():
    print("Job is still running...")
    time.sleep(5)

print(f"Job is complete! Status: {sync_result.get_job_status()}")

Examining the sync result

You can examine the sync result to get more information about the job:

sync_result: SyncResult = connection.run_sync()

# Print the job details
print(
    f'''
    Job ID: {sync_result.job_id}
    Job URL: {sync_result.job_url}
    Start Time: {sync_result.start_time}
    Records Synced: {sync_result.records_synced}
    Bytes Synced: {sync_result.bytes_synced}
    Job Status: {sync_result.get_job_status()}
    List of Stream Names: {', '.join(sync_result.stream_names)}
    '''
)

Reading data from Airbyte Cloud sync result

This feature is currently only available for specific SQL-based destinations. This includes SQL-based destinations such as Snowflake and BigQuery. The list of supported destinations may be determined by inspecting the constant airbyte.cloud.constants.READABLE_DESTINATION_TYPES.

If your destination is supported, you can read records directly from the SyncResult object.

# Assuming we've already created a `connection` object...
sync_result = connection.get_sync_result()

# Print a list of available stream names
print(sync_result.stream_names)

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

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

# Or iterate over the dataset directly
for record in dataset:
    print(record)

  1# Copyright (c) 2024 Airbyte, Inc., all rights reserved.
  2"""Sync results for Airbyte Cloud workspaces.
  3
  4## Examples
  5
  6### Run a sync job and wait for completion
  7
  8To get started, we'll need a `.CloudConnection` object. You can obtain this object by calling
  9`.CloudWorkspace.get_connection()`.
 10
 11```python
 12from airbyte import cloud
 13
 14# Initialize an Airbyte Cloud workspace object
 15workspace = cloud.CloudWorkspace(
 16    workspace_id="123",
 17    api_key=ab.get_secret("AIRBYTE_CLOUD_API_KEY"),
 18)
 19
 20# Get a connection object
 21connection = workspace.get_connection(connection_id="456")
 22```
 23
 24Once we have a `.CloudConnection` object, we can simply call `run_sync()`
 25to start a sync job and wait for it to complete.
 26
 27```python
 28# Run a sync job
 29sync_result: SyncResult = connection.run_sync()
 30```
 31
 32### Run a sync job and return immediately
 33
 34By default, `run_sync()` will wait for the job to complete and raise an
 35exception if the job fails. You can instead return immediately by setting
 36`wait=False`.
 37
 38```python
 39# Start the sync job and return immediately
 40sync_result: SyncResult = connection.run_sync(wait=False)
 41
 42while not sync_result.is_job_complete():
 43    print("Job is still running...")
 44    time.sleep(5)
 45
 46print(f"Job is complete! Status: {sync_result.get_job_status()}")
 47```
 48
 49### Examining the sync result
 50
 51You can examine the sync result to get more information about the job:
 52
 53```python
 54sync_result: SyncResult = connection.run_sync()
 55
 56# Print the job details
 57print(
 58    f'''
 59    Job ID: {sync_result.job_id}
 60    Job URL: {sync_result.job_url}
 61    Start Time: {sync_result.start_time}
 62    Records Synced: {sync_result.records_synced}
 63    Bytes Synced: {sync_result.bytes_synced}
 64    Job Status: {sync_result.get_job_status()}
 65    List of Stream Names: {', '.join(sync_result.stream_names)}
 66    '''
 67)
 68```
 69
 70### Reading data from Airbyte Cloud sync result
 71
 72**This feature is currently only available for specific SQL-based destinations.** This includes
 73SQL-based destinations such as Snowflake and BigQuery. The list of supported destinations may be
 74determined by inspecting the constant `airbyte.cloud.constants.READABLE_DESTINATION_TYPES`.
 75
 76If your destination is supported, you can read records directly from the SyncResult object.
 77
 78```python
 79# Assuming we've already created a `connection` object...
 80sync_result = connection.get_sync_result()
 81
 82# Print a list of available stream names
 83print(sync_result.stream_names)
 84
 85# Get a dataset from the sync result
 86dataset: CachedDataset = sync_result.get_dataset("users")
 87
 88# Get the SQLAlchemy table to use in SQL queries...
 89users_table = dataset.to_sql_table()
 90print(f"Table name: {users_table.name}")
 91
 92# Or iterate over the dataset directly
 93for record in dataset:
 94    print(record)
 95```
 96
 97------
 98
 99"""
100
101from __future__ import annotations
102
103import time
104from collections.abc import Iterator, Mapping
105from dataclasses import asdict, dataclass
106from typing import TYPE_CHECKING, Any
107
108from typing_extensions import final
109
110from airbyte_cdk.utils.datetime_helpers import ab_datetime_parse
111
112from airbyte._util import api_util
113from airbyte.caches._utils._dest_to_cache import destination_to_cache
114from airbyte.cloud.constants import FAILED_STATUSES, FINAL_STATUSES
115from airbyte.cloud.models import CloudConnectionInfo, CloudJobInfo, JobStatusEnum
116from airbyte.datasets import CachedDataset
117from airbyte.exceptions import AirbyteConnectionSyncError, AirbyteConnectionSyncTimeoutError
118
119
120DEFAULT_SYNC_TIMEOUT_SECONDS = 30 * 60  # 30 minutes
121"""The default timeout for waiting for a sync job to complete, in seconds."""
122
123if TYPE_CHECKING:
124    from datetime import datetime
125
126    import sqlalchemy
127
128    from airbyte.caches.base import CacheBase
129    from airbyte.cloud.connections import CloudConnection
130    from airbyte.cloud.workspaces import CloudWorkspace
131
132
133@dataclass
134class SyncAttempt:
135    """Represents a single attempt of a sync job.
136
137    **This class is not meant to be instantiated directly.** Instead, obtain a `SyncAttempt` by
138    calling `.SyncResult.get_attempts()`.
139    """
140
141    workspace: CloudWorkspace
142    connection: CloudConnection
143    job_id: int
144    attempt_number: int
145    _attempt_data: dict[str, Any] | None = None
146
147    @property
148    def attempt_id(self) -> int:
149        """Return the attempt ID."""
150        return self._get_attempt_data()["id"]
151
152    @property
153    def status(self) -> str:
154        """Return the attempt status."""
155        return self._get_attempt_data()["status"]
156
157    @property
158    def bytes_synced(self) -> int:
159        """Return the number of bytes synced in this attempt."""
160        return self._get_attempt_data().get("bytesSynced", 0)
161
162    @property
163    def records_synced(self) -> int:
164        """Return the number of records synced in this attempt."""
165        return self._get_attempt_data().get("recordsSynced", 0)
166
167    @property
168    def created_at(self) -> datetime:
169        """Return the creation time of the attempt."""
170        timestamp = self._get_attempt_data()["createdAt"]
171        return ab_datetime_parse(timestamp)
172
173    def _get_attempt_data(self) -> dict[str, Any]:
174        """Get attempt data from the provided attempt data."""
175        if self._attempt_data is None:
176            raise ValueError(
177                "Attempt data not provided. SyncAttempt should be created via "
178                "SyncResult.get_attempts()."
179            )
180        return self._attempt_data["attempt"]
181
182    def get_full_log_text(self) -> str:
183        """Return the complete log text for this attempt.
184
185        Returns:
186            String containing all log text for this attempt, with lines separated by newlines.
187        """
188        if self._attempt_data is None:
189            return ""
190
191        logs_data = self._attempt_data.get("logs")
192        if not logs_data:
193            return ""
194
195        result = ""
196
197        if "events" in logs_data:
198            log_events = logs_data["events"]
199            if log_events:
200                log_lines = []
201                for event in log_events:
202                    timestamp = event.get("timestamp", "")
203                    level = event.get("level", "INFO")
204                    message = event.get("message", "")
205                    log_lines.append(
206                        f"[{timestamp}] {level}: {message}"  # pyrefly: ignore[bad-argument-type]
207                    )
208                result = "\n".join(log_lines)
209        elif "logLines" in logs_data:
210            log_lines = logs_data["logLines"]
211            if log_lines:
212                result = "\n".join(log_lines)
213
214        return result
215
216
217@dataclass
218class SyncResult:
219    """The result of a sync operation.
220
221    **This class is not meant to be instantiated directly.** Instead, obtain a `SyncResult` by
222    interacting with the `.CloudWorkspace` and `.CloudConnection` objects.
223    """
224
225    workspace: CloudWorkspace
226    connection: CloudConnection
227    job_id: int
228    table_name_prefix: str = ""
229    table_name_suffix: str = ""
230    _latest_job_info: CloudJobInfo | None = None
231    _connection_response: CloudConnectionInfo | None = None
232    _cache: CacheBase | None = None
233    _job_with_attempts_info: dict[str, Any] | None = None
234
235    @property
236    def job_url(self) -> str:
237        """Return the URL of the sync job.
238
239        Note: This currently returns the connection's job history URL, as there is no direct URL
240        to a specific job in the Airbyte Cloud web app.
241
242        TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number.
243              E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true
244        """
245        return f"{self.connection.job_history_url}"
246
247    def _get_connection_info(self, *, force_refresh: bool = False) -> CloudConnectionInfo:
248        """Return connection info for the sync job."""
249        if self._connection_response and not force_refresh:
250            return self._connection_response
251
252        self._connection_response = CloudConnectionInfo.from_api_response(
253            api_util.get_connection(
254                workspace_id=self.workspace.workspace_id,
255                api_root=self.workspace.api_root,
256                connection_id=self.connection.connection_id,
257                client_id=self.workspace.client_id,
258                client_secret=self.workspace.client_secret,
259                bearer_token=self.workspace.bearer_token,
260            )
261        )
262        return self._connection_response
263
264    def _get_destination_configuration(self, *, force_refresh: bool = False) -> dict[str, Any]:
265        """Return the destination configuration for the sync job."""
266        connection_info = self._get_connection_info(force_refresh=force_refresh)
267        destination_response = api_util.get_destination(
268            destination_id=connection_info.destination_id,
269            api_root=self.workspace.api_root,
270            client_id=self.workspace.client_id,
271            client_secret=self.workspace.client_secret,
272            bearer_token=self.workspace.bearer_token,
273        )
274        configuration = destination_response.configuration
275        if isinstance(configuration, Mapping):
276            configuration_dict = configuration
277        else:
278            configuration_dict = asdict(configuration)
279        return {**configuration_dict, "destinationType": destination_response.destination_type}
280
281    def is_job_complete(self) -> bool:
282        """Check if the sync job is complete."""
283        return self.get_job_status() in FINAL_STATUSES
284
285    def get_job_status(self) -> JobStatusEnum:
286        """Check if the sync job is still running."""
287        return self._fetch_latest_job_info().status
288
289    def _fetch_latest_job_info(self) -> CloudJobInfo:
290        """Return the job info for the sync job."""
291        if self._latest_job_info and self._latest_job_info.status in FINAL_STATUSES:
292            return self._latest_job_info
293
294        self._latest_job_info = CloudJobInfo.from_api_response(
295            api_util.get_job_info(
296                job_id=self.job_id,
297                api_root=self.workspace.api_root,
298                client_id=self.workspace.client_id,
299                client_secret=self.workspace.client_secret,
300                bearer_token=self.workspace.bearer_token,
301            )
302        )
303        return self._latest_job_info
304
305    @property
306    def bytes_synced(self) -> int:
307        """Return the number of records processed."""
308        return self._fetch_latest_job_info().bytes_synced or 0
309
310    @property
311    def records_synced(self) -> int:
312        """Return the number of records processed."""
313        return self._fetch_latest_job_info().rows_synced or 0
314
315    @property
316    def start_time(self) -> datetime:
317        """Return the start time of the sync job in UTC."""
318        try:
319            return ab_datetime_parse(self._fetch_latest_job_info().start_time)
320        except (ValueError, TypeError) as e:
321            if "Invalid isoformat string" in str(e):
322                job_info_raw = api_util._make_config_api_request(  # noqa: SLF001
323                    api_root=self.workspace.api_root,
324                    config_api_root=self.workspace.config_api_root,
325                    path="/jobs/get",
326                    json={"id": self.job_id},
327                    client_id=self.workspace.client_id,
328                    client_secret=self.workspace.client_secret,
329                    bearer_token=self.workspace.bearer_token,
330                )
331                raw_start_time = job_info_raw.get("startTime")
332                if raw_start_time:
333                    return ab_datetime_parse(raw_start_time)
334            raise
335
336    def _fetch_job_with_attempts(self) -> dict[str, Any]:
337        """Fetch job info with attempts from Config API using lazy loading pattern."""
338        if self._job_with_attempts_info is not None:
339            return self._job_with_attempts_info
340
341        self._job_with_attempts_info = api_util._make_config_api_request(  # noqa: SLF001  # Config API helper
342            api_root=self.workspace.api_root,
343            config_api_root=self.workspace.config_api_root,
344            path="/jobs/get",
345            json={
346                "id": self.job_id,
347            },
348            client_id=self.workspace.client_id,
349            client_secret=self.workspace.client_secret,
350            bearer_token=self.workspace.bearer_token,
351        )
352        return self._job_with_attempts_info
353
354    def get_attempts(self) -> list[SyncAttempt]:
355        """Return a list of attempts for this sync job."""
356        job_with_attempts = self._fetch_job_with_attempts()
357        attempts_data = job_with_attempts.get("attempts", [])
358
359        return [
360            SyncAttempt(
361                workspace=self.workspace,
362                connection=self.connection,
363                job_id=self.job_id,
364                attempt_number=i,
365                _attempt_data=attempt_data,
366            )
367            for i, attempt_data in enumerate(attempts_data, start=0)
368        ]
369
370    def raise_failure_status(
371        self,
372        *,
373        refresh_status: bool = False,
374    ) -> None:
375        """Raise an exception if the sync job failed.
376
377        By default, this method will use the latest status available. If you want to refresh the
378        status before checking for failure, set `refresh_status=True`. If the job has failed, this
379        method will raise a `AirbyteConnectionSyncError`.
380
381        Otherwise, do nothing.
382        """
383        if not refresh_status and self._latest_job_info:
384            latest_status = self._latest_job_info.status
385        else:
386            latest_status = self.get_job_status()
387
388        if latest_status in FAILED_STATUSES:
389            raise AirbyteConnectionSyncError(
390                workspace=self.workspace,
391                connection_id=self.connection.connection_id,
392                job_id=self.job_id,
393                job_status=self.get_job_status(),
394            )
395
396    def wait_for_completion(
397        self,
398        *,
399        wait_timeout: int = DEFAULT_SYNC_TIMEOUT_SECONDS,
400        raise_timeout: bool = True,
401        raise_failure: bool = False,
402    ) -> JobStatusEnum:
403        """Wait for a job to finish running."""
404        start_time = time.time()
405        while True:
406            latest_status = self.get_job_status()
407            if latest_status in FINAL_STATUSES:
408                if raise_failure:
409                    # No-op if the job succeeded or is still running:
410                    self.raise_failure_status()
411
412                return latest_status
413
414            if time.time() - start_time > wait_timeout:
415                if raise_timeout:
416                    raise AirbyteConnectionSyncTimeoutError(
417                        workspace=self.workspace,
418                        connection_id=self.connection.connection_id,
419                        job_id=self.job_id,
420                        job_status=latest_status,
421                        timeout=wait_timeout,
422                    )
423
424                return latest_status  # This will be a non-final status
425
426            time.sleep(api_util.JOB_WAIT_INTERVAL_SECS)
427
428    def get_sql_cache(self) -> CacheBase:
429        """Return a SQL Cache object for working with the data in a SQL-based destination's."""
430        if self._cache:
431            return self._cache
432
433        destination_configuration = self._get_destination_configuration()
434        self._cache = destination_to_cache(destination_configuration=destination_configuration)
435        return self._cache
436
437    def get_sql_engine(self) -> sqlalchemy.engine.Engine:
438        """Return a SQL Engine for querying a SQL-based destination."""
439        return self.get_sql_cache().get_sql_engine()
440
441    def get_sql_table_name(self, stream_name: str) -> str:
442        """Return the SQL table name of the named stream."""
443        return self.get_sql_cache().processor.get_sql_table_name(stream_name=stream_name)
444
445    def get_sql_table(
446        self,
447        stream_name: str,
448    ) -> sqlalchemy.Table:
449        """Return a SQLAlchemy table object for the named stream."""
450        return self.get_sql_cache().processor.get_sql_table(stream_name)
451
452    def get_dataset(self, stream_name: str) -> CachedDataset:
453        """Retrieve an `airbyte.datasets.CachedDataset` object for a given stream name.
454
455        This can be used to read and analyze the data in a SQL-based destination.
456
457        TODO: In a future iteration, we can consider providing stream configuration information
458              (catalog information) to the `CachedDataset` object via the "Get stream properties"
459              API: https://reference.airbyte.com/reference/getstreamproperties
460        """
461        return CachedDataset(
462            self.get_sql_cache(),
463            stream_name=stream_name,
464            stream_configuration=False,  # Don't look for stream configuration in cache.
465        )
466
467    def get_sql_database_name(self) -> str:
468        """Return the SQL database name."""
469        cache = self.get_sql_cache()
470        return cache.get_database_name()
471
472    def get_sql_schema_name(self) -> str:
473        """Return the SQL schema name."""
474        cache = self.get_sql_cache()
475        return cache.schema_name
476
477    @property
478    def stream_names(self) -> list[str]:
479        """Return the set of stream names."""
480        return self.connection.stream_names
481
482    @final
483    @property
484    def streams(
485        self,
486    ) -> _SyncResultStreams:  # pyrefly: ignore[unknown-name]
487        """Return a mapping of stream names to `airbyte.CachedDataset` objects.
488
489        This is a convenience wrapper around the `stream_names`
490        property and `get_dataset()` method.
491        """
492        return self._SyncResultStreams(self)
493
494    class _SyncResultStreams(Mapping[str, CachedDataset]):
495        """A mapping of stream names to cached datasets."""
496
497        def __init__(
498            self,
499            parent: SyncResult,
500            /,
501        ) -> None:
502            self.parent: SyncResult = parent
503
504        def __getitem__(self, key: str) -> CachedDataset:
505            return self.parent.get_dataset(stream_name=key)
506
507        def __iter__(self) -> Iterator[str]:
508            return iter(self.parent.stream_names)
509
510        def __len__(self) -> int:
511            return len(self.parent.stream_names)
512
513
514__all__ = [
515    "SyncResult",
516    "SyncAttempt",
517]
@dataclass
class SyncResult:
218@dataclass
219class SyncResult:
220    """The result of a sync operation.
221
222    **This class is not meant to be instantiated directly.** Instead, obtain a `SyncResult` by
223    interacting with the `.CloudWorkspace` and `.CloudConnection` objects.
224    """
225
226    workspace: CloudWorkspace
227    connection: CloudConnection
228    job_id: int
229    table_name_prefix: str = ""
230    table_name_suffix: str = ""
231    _latest_job_info: CloudJobInfo | None = None
232    _connection_response: CloudConnectionInfo | None = None
233    _cache: CacheBase | None = None
234    _job_with_attempts_info: dict[str, Any] | None = None
235
236    @property
237    def job_url(self) -> str:
238        """Return the URL of the sync job.
239
240        Note: This currently returns the connection's job history URL, as there is no direct URL
241        to a specific job in the Airbyte Cloud web app.
242
243        TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number.
244              E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true
245        """
246        return f"{self.connection.job_history_url}"
247
248    def _get_connection_info(self, *, force_refresh: bool = False) -> CloudConnectionInfo:
249        """Return connection info for the sync job."""
250        if self._connection_response and not force_refresh:
251            return self._connection_response
252
253        self._connection_response = CloudConnectionInfo.from_api_response(
254            api_util.get_connection(
255                workspace_id=self.workspace.workspace_id,
256                api_root=self.workspace.api_root,
257                connection_id=self.connection.connection_id,
258                client_id=self.workspace.client_id,
259                client_secret=self.workspace.client_secret,
260                bearer_token=self.workspace.bearer_token,
261            )
262        )
263        return self._connection_response
264
265    def _get_destination_configuration(self, *, force_refresh: bool = False) -> dict[str, Any]:
266        """Return the destination configuration for the sync job."""
267        connection_info = self._get_connection_info(force_refresh=force_refresh)
268        destination_response = api_util.get_destination(
269            destination_id=connection_info.destination_id,
270            api_root=self.workspace.api_root,
271            client_id=self.workspace.client_id,
272            client_secret=self.workspace.client_secret,
273            bearer_token=self.workspace.bearer_token,
274        )
275        configuration = destination_response.configuration
276        if isinstance(configuration, Mapping):
277            configuration_dict = configuration
278        else:
279            configuration_dict = asdict(configuration)
280        return {**configuration_dict, "destinationType": destination_response.destination_type}
281
282    def is_job_complete(self) -> bool:
283        """Check if the sync job is complete."""
284        return self.get_job_status() in FINAL_STATUSES
285
286    def get_job_status(self) -> JobStatusEnum:
287        """Check if the sync job is still running."""
288        return self._fetch_latest_job_info().status
289
290    def _fetch_latest_job_info(self) -> CloudJobInfo:
291        """Return the job info for the sync job."""
292        if self._latest_job_info and self._latest_job_info.status in FINAL_STATUSES:
293            return self._latest_job_info
294
295        self._latest_job_info = CloudJobInfo.from_api_response(
296            api_util.get_job_info(
297                job_id=self.job_id,
298                api_root=self.workspace.api_root,
299                client_id=self.workspace.client_id,
300                client_secret=self.workspace.client_secret,
301                bearer_token=self.workspace.bearer_token,
302            )
303        )
304        return self._latest_job_info
305
306    @property
307    def bytes_synced(self) -> int:
308        """Return the number of records processed."""
309        return self._fetch_latest_job_info().bytes_synced or 0
310
311    @property
312    def records_synced(self) -> int:
313        """Return the number of records processed."""
314        return self._fetch_latest_job_info().rows_synced or 0
315
316    @property
317    def start_time(self) -> datetime:
318        """Return the start time of the sync job in UTC."""
319        try:
320            return ab_datetime_parse(self._fetch_latest_job_info().start_time)
321        except (ValueError, TypeError) as e:
322            if "Invalid isoformat string" in str(e):
323                job_info_raw = api_util._make_config_api_request(  # noqa: SLF001
324                    api_root=self.workspace.api_root,
325                    config_api_root=self.workspace.config_api_root,
326                    path="/jobs/get",
327                    json={"id": self.job_id},
328                    client_id=self.workspace.client_id,
329                    client_secret=self.workspace.client_secret,
330                    bearer_token=self.workspace.bearer_token,
331                )
332                raw_start_time = job_info_raw.get("startTime")
333                if raw_start_time:
334                    return ab_datetime_parse(raw_start_time)
335            raise
336
337    def _fetch_job_with_attempts(self) -> dict[str, Any]:
338        """Fetch job info with attempts from Config API using lazy loading pattern."""
339        if self._job_with_attempts_info is not None:
340            return self._job_with_attempts_info
341
342        self._job_with_attempts_info = api_util._make_config_api_request(  # noqa: SLF001  # Config API helper
343            api_root=self.workspace.api_root,
344            config_api_root=self.workspace.config_api_root,
345            path="/jobs/get",
346            json={
347                "id": self.job_id,
348            },
349            client_id=self.workspace.client_id,
350            client_secret=self.workspace.client_secret,
351            bearer_token=self.workspace.bearer_token,
352        )
353        return self._job_with_attempts_info
354
355    def get_attempts(self) -> list[SyncAttempt]:
356        """Return a list of attempts for this sync job."""
357        job_with_attempts = self._fetch_job_with_attempts()
358        attempts_data = job_with_attempts.get("attempts", [])
359
360        return [
361            SyncAttempt(
362                workspace=self.workspace,
363                connection=self.connection,
364                job_id=self.job_id,
365                attempt_number=i,
366                _attempt_data=attempt_data,
367            )
368            for i, attempt_data in enumerate(attempts_data, start=0)
369        ]
370
371    def raise_failure_status(
372        self,
373        *,
374        refresh_status: bool = False,
375    ) -> None:
376        """Raise an exception if the sync job failed.
377
378        By default, this method will use the latest status available. If you want to refresh the
379        status before checking for failure, set `refresh_status=True`. If the job has failed, this
380        method will raise a `AirbyteConnectionSyncError`.
381
382        Otherwise, do nothing.
383        """
384        if not refresh_status and self._latest_job_info:
385            latest_status = self._latest_job_info.status
386        else:
387            latest_status = self.get_job_status()
388
389        if latest_status in FAILED_STATUSES:
390            raise AirbyteConnectionSyncError(
391                workspace=self.workspace,
392                connection_id=self.connection.connection_id,
393                job_id=self.job_id,
394                job_status=self.get_job_status(),
395            )
396
397    def wait_for_completion(
398        self,
399        *,
400        wait_timeout: int = DEFAULT_SYNC_TIMEOUT_SECONDS,
401        raise_timeout: bool = True,
402        raise_failure: bool = False,
403    ) -> JobStatusEnum:
404        """Wait for a job to finish running."""
405        start_time = time.time()
406        while True:
407            latest_status = self.get_job_status()
408            if latest_status in FINAL_STATUSES:
409                if raise_failure:
410                    # No-op if the job succeeded or is still running:
411                    self.raise_failure_status()
412
413                return latest_status
414
415            if time.time() - start_time > wait_timeout:
416                if raise_timeout:
417                    raise AirbyteConnectionSyncTimeoutError(
418                        workspace=self.workspace,
419                        connection_id=self.connection.connection_id,
420                        job_id=self.job_id,
421                        job_status=latest_status,
422                        timeout=wait_timeout,
423                    )
424
425                return latest_status  # This will be a non-final status
426
427            time.sleep(api_util.JOB_WAIT_INTERVAL_SECS)
428
429    def get_sql_cache(self) -> CacheBase:
430        """Return a SQL Cache object for working with the data in a SQL-based destination's."""
431        if self._cache:
432            return self._cache
433
434        destination_configuration = self._get_destination_configuration()
435        self._cache = destination_to_cache(destination_configuration=destination_configuration)
436        return self._cache
437
438    def get_sql_engine(self) -> sqlalchemy.engine.Engine:
439        """Return a SQL Engine for querying a SQL-based destination."""
440        return self.get_sql_cache().get_sql_engine()
441
442    def get_sql_table_name(self, stream_name: str) -> str:
443        """Return the SQL table name of the named stream."""
444        return self.get_sql_cache().processor.get_sql_table_name(stream_name=stream_name)
445
446    def get_sql_table(
447        self,
448        stream_name: str,
449    ) -> sqlalchemy.Table:
450        """Return a SQLAlchemy table object for the named stream."""
451        return self.get_sql_cache().processor.get_sql_table(stream_name)
452
453    def get_dataset(self, stream_name: str) -> CachedDataset:
454        """Retrieve an `airbyte.datasets.CachedDataset` object for a given stream name.
455
456        This can be used to read and analyze the data in a SQL-based destination.
457
458        TODO: In a future iteration, we can consider providing stream configuration information
459              (catalog information) to the `CachedDataset` object via the "Get stream properties"
460              API: https://reference.airbyte.com/reference/getstreamproperties
461        """
462        return CachedDataset(
463            self.get_sql_cache(),
464            stream_name=stream_name,
465            stream_configuration=False,  # Don't look for stream configuration in cache.
466        )
467
468    def get_sql_database_name(self) -> str:
469        """Return the SQL database name."""
470        cache = self.get_sql_cache()
471        return cache.get_database_name()
472
473    def get_sql_schema_name(self) -> str:
474        """Return the SQL schema name."""
475        cache = self.get_sql_cache()
476        return cache.schema_name
477
478    @property
479    def stream_names(self) -> list[str]:
480        """Return the set of stream names."""
481        return self.connection.stream_names
482
483    @final
484    @property
485    def streams(
486        self,
487    ) -> _SyncResultStreams:  # pyrefly: ignore[unknown-name]
488        """Return a mapping of stream names to `airbyte.CachedDataset` objects.
489
490        This is a convenience wrapper around the `stream_names`
491        property and `get_dataset()` method.
492        """
493        return self._SyncResultStreams(self)
494
495    class _SyncResultStreams(Mapping[str, CachedDataset]):
496        """A mapping of stream names to cached datasets."""
497
498        def __init__(
499            self,
500            parent: SyncResult,
501            /,
502        ) -> None:
503            self.parent: SyncResult = parent
504
505        def __getitem__(self, key: str) -> CachedDataset:
506            return self.parent.get_dataset(stream_name=key)
507
508        def __iter__(self) -> Iterator[str]:
509            return iter(self.parent.stream_names)
510
511        def __len__(self) -> int:
512            return len(self.parent.stream_names)

The result of a sync operation.

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

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

Return the URL of the sync job.

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

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

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

Check if the sync job is complete.

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

Check if the sync job is still running.

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

Return the number of records processed.

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

Return the number of records processed.

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

Return the start time of the sync job in UTC.

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

Return a list of attempts for this sync job.

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

Raise an exception if the sync job failed.

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

Otherwise, do nothing.

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

Wait for a job to finish running.

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

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

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

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

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

Return the SQL table name of the named stream.

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

Return a SQLAlchemy table object for the named stream.

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

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

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

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

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

Return the SQL database name.

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

Return the SQL schema name.

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

Return the set of stream names.

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

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

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

@dataclass
class SyncAttempt:
134@dataclass
135class SyncAttempt:
136    """Represents a single attempt of a sync job.
137
138    **This class is not meant to be instantiated directly.** Instead, obtain a `SyncAttempt` by
139    calling `.SyncResult.get_attempts()`.
140    """
141
142    workspace: CloudWorkspace
143    connection: CloudConnection
144    job_id: int
145    attempt_number: int
146    _attempt_data: dict[str, Any] | None = None
147
148    @property
149    def attempt_id(self) -> int:
150        """Return the attempt ID."""
151        return self._get_attempt_data()["id"]
152
153    @property
154    def status(self) -> str:
155        """Return the attempt status."""
156        return self._get_attempt_data()["status"]
157
158    @property
159    def bytes_synced(self) -> int:
160        """Return the number of bytes synced in this attempt."""
161        return self._get_attempt_data().get("bytesSynced", 0)
162
163    @property
164    def records_synced(self) -> int:
165        """Return the number of records synced in this attempt."""
166        return self._get_attempt_data().get("recordsSynced", 0)
167
168    @property
169    def created_at(self) -> datetime:
170        """Return the creation time of the attempt."""
171        timestamp = self._get_attempt_data()["createdAt"]
172        return ab_datetime_parse(timestamp)
173
174    def _get_attempt_data(self) -> dict[str, Any]:
175        """Get attempt data from the provided attempt data."""
176        if self._attempt_data is None:
177            raise ValueError(
178                "Attempt data not provided. SyncAttempt should be created via "
179                "SyncResult.get_attempts()."
180            )
181        return self._attempt_data["attempt"]
182
183    def get_full_log_text(self) -> str:
184        """Return the complete log text for this attempt.
185
186        Returns:
187            String containing all log text for this attempt, with lines separated by newlines.
188        """
189        if self._attempt_data is None:
190            return ""
191
192        logs_data = self._attempt_data.get("logs")
193        if not logs_data:
194            return ""
195
196        result = ""
197
198        if "events" in logs_data:
199            log_events = logs_data["events"]
200            if log_events:
201                log_lines = []
202                for event in log_events:
203                    timestamp = event.get("timestamp", "")
204                    level = event.get("level", "INFO")
205                    message = event.get("message", "")
206                    log_lines.append(
207                        f"[{timestamp}] {level}: {message}"  # pyrefly: ignore[bad-argument-type]
208                    )
209                result = "\n".join(log_lines)
210        elif "logLines" in logs_data:
211            log_lines = logs_data["logLines"]
212            if log_lines:
213                result = "\n".join(log_lines)
214
215        return result

Represents a single attempt of a sync job.

This class is not meant to be instantiated directly. Instead, obtain a SyncAttempt by calling .SyncResult.get_attempts().

SyncAttempt( workspace: airbyte.cloud.CloudWorkspace, connection: airbyte.cloud.CloudConnection, job_id: int, attempt_number: int, _attempt_data: dict[str, typing.Any] | None = None)
job_id: int
attempt_number: int
attempt_id: int
148    @property
149    def attempt_id(self) -> int:
150        """Return the attempt ID."""
151        return self._get_attempt_data()["id"]

Return the attempt ID.

status: str
153    @property
154    def status(self) -> str:
155        """Return the attempt status."""
156        return self._get_attempt_data()["status"]

Return the attempt status.

bytes_synced: int
158    @property
159    def bytes_synced(self) -> int:
160        """Return the number of bytes synced in this attempt."""
161        return self._get_attempt_data().get("bytesSynced", 0)

Return the number of bytes synced in this attempt.

records_synced: int
163    @property
164    def records_synced(self) -> int:
165        """Return the number of records synced in this attempt."""
166        return self._get_attempt_data().get("recordsSynced", 0)

Return the number of records synced in this attempt.

created_at: datetime.datetime
168    @property
169    def created_at(self) -> datetime:
170        """Return the creation time of the attempt."""
171        timestamp = self._get_attempt_data()["createdAt"]
172        return ab_datetime_parse(timestamp)

Return the creation time of the attempt.

def get_full_log_text(self) -> str:
183    def get_full_log_text(self) -> str:
184        """Return the complete log text for this attempt.
185
186        Returns:
187            String containing all log text for this attempt, with lines separated by newlines.
188        """
189        if self._attempt_data is None:
190            return ""
191
192        logs_data = self._attempt_data.get("logs")
193        if not logs_data:
194            return ""
195
196        result = ""
197
198        if "events" in logs_data:
199            log_events = logs_data["events"]
200            if log_events:
201                log_lines = []
202                for event in log_events:
203                    timestamp = event.get("timestamp", "")
204                    level = event.get("level", "INFO")
205                    message = event.get("message", "")
206                    log_lines.append(
207                        f"[{timestamp}] {level}: {message}"  # pyrefly: ignore[bad-argument-type]
208                    )
209                result = "\n".join(log_lines)
210        elif "logLines" in logs_data:
211            log_lines = logs_data["logLines"]
212            if log_lines:
213                result = "\n".join(log_lines)
214
215        return result

Return the complete log text for this attempt.

Returns:

String containing all log text for this attempt, with lines separated by newlines.