From 20eccfa5e9542a6928b90956317730bc9d09f562 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 12:03:36 -0600 Subject: [PATCH 01/17] feat: Attempt to sync all streams instead of crashing on the first error MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/__init__.py | 4 +- singer_sdk/streams/core.py | 35 +++ singer_sdk/tap_base.py | 130 +++++++++-- .../test_child_deselected_parent/stderr.log | 2 + .../test_deselected_child/stderr.log | 1 + .../test_one_parent_many_children/stderr.log | 2 + .../stderr.log | 2 + .../stderr.log | 2 + tests/core/test_sync_outcomes.py | 214 ++++++++++++++++++ .../test_countries_to_csv/singer.log | 2 + .../activate_version/singer.log | 2 + .../no_activate_version/singer.log | 2 + .../test_fake_people_to_csv/singer.log | 1 + 13 files changed, 373 insertions(+), 26 deletions(-) create mode 100644 tests/core/test_sync_outcomes.py diff --git a/singer_sdk/streams/__init__.py b/singer_sdk/streams/__init__.py index f95dc12c2b..a053233451 100644 --- a/singer_sdk/streams/__init__.py +++ b/singer_sdk/streams/__init__.py @@ -6,11 +6,11 @@ import warnings from singer_sdk.helpers._compat import SingerSDKDeprecationWarning -from singer_sdk.streams.core import Stream +from singer_sdk.streams.core import Stream, SyncResult from singer_sdk.streams.graphql import GraphQLStream from singer_sdk.streams.rest import RESTStream -__all__ = ["GraphQLStream", "RESTStream", "Stream"] +__all__ = ["GraphQLStream", "RESTStream", "Stream", "SyncResult"] def __getattr__(name: str) -> t.Any: # noqa: ANN401 diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index 631faab117..ebeee28ce4 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -5,6 +5,7 @@ import abc import copy import datetime +import enum import json import logging import typing as t @@ -54,6 +55,27 @@ from singer_sdk.tap_base import Tap +class SyncResult(enum.Enum): + """Outcome of a single stream's sync operation. + + Set on :attr:`~singer_sdk.Stream.sync_result` after + :meth:`~singer_sdk.Stream.sync` completes or fails. + + Attributes: + SUCCESS: Completed without error. + FAILED: Raised a fatal (non-lifecycle) exception. + ABORTED: Raised a lifecycle abort exception + (:class:`~singer_sdk.exceptions.AbortedSyncFailedException` or + :class:`~singer_sdk.exceptions.AbortedSyncPausedException`). + PARTIAL: Reserved — ignorable errors with skipped records (requires PR 3). + """ + + SUCCESS = "success" + FAILED = "failed" + ABORTED = "aborted" + PARTIAL = "partial" + + class Stream(abc.ABC): # noqa: PLR0904 """Abstract base class for tap streams. @@ -155,6 +177,7 @@ def __init__( self._schema: dict | None = None self._sync_costs: dict[str, int] = {} self.child_streams: list[Stream] = [] + self.sync_result: SyncResult | None = None # Initialize state manager self._state_manager = StreamStateManager( @@ -1311,6 +1334,10 @@ def sync(self, context: types.Context | None = None) -> None: Args: context: Stream partition or context dictionary. + + Raises: + AbortedSyncFailedException: If the sync was aborted non-resumably. + AbortedSyncPausedException: If the sync was paused at a resumable point. """ # Preprocess context before it's frozen context = self.preprocess_context(context) if context else None @@ -1349,13 +1376,21 @@ def sync(self, context: types.Context | None = None) -> None: # Sync the records themselves: for _ in self._sync_records(context=context): pass + except (AbortedSyncFailedException, AbortedSyncPausedException): + # Lifecycle abort — must not be swallowed. Mark ABORTED and re-raise + # so sync_all() can propagate to invoke() for proper exit-code handling. + self.sync_result = SyncResult.ABORTED + raise except Exception: self.log( "An unhandled error occurred while syncing '%s'", self.name, level=logging.ERROR, ) + self.sync_result = SyncResult.FAILED raise + else: + self.sync_result = SyncResult.SUCCESS def _sync_children(self, child_context: types.Context | None) -> None: if child_context is None: diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index 3d2f454bb1..5fd27e5753 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -5,6 +5,8 @@ import abc import collections.abc import contextlib +import logging +import sys import typing as t import warnings from enum import Enum @@ -29,6 +31,7 @@ from singer_sdk.io_base import SingerWriter from singer_sdk.plugin_base import BaseSingerWriter, PluginBase, _ConfigInput from singer_sdk.singerlib import Catalog +from singer_sdk.streams.core import SyncResult if t.TYPE_CHECKING: from pathlib import PurePath @@ -469,37 +472,95 @@ def _set_compatible_replication_methods(self) -> None: # Sync methods + def _log_stream_sync_result(self, stream: Stream) -> None: + """Log a one-line sync outcome for a stream. + + Called once per stream at the end of :meth:`sync_all`. + + Args: + stream: The stream whose result to log. + """ + result = stream.sync_result + if result is None: + # Stream was never directly synced (deselected or child stream). + return + + level_map = { + SyncResult.SUCCESS: logging.INFO, + SyncResult.FAILED: logging.ERROR, + SyncResult.ABORTED: logging.WARNING, + SyncResult.PARTIAL: logging.WARNING, + } + self.logger.log( + level_map.get(result, logging.INFO), + "Stream '%s' sync result: %s", + stream.name, + result.value, + ) + @t.final def sync_all(self) -> None: - """Sync all streams.""" + """Sync all streams. + + A stream that raises a fatal (non-lifecycle) exception is logged and + skipped; syncing continues with remaining streams. + + :class:`~singer_sdk.exceptions.AbortedSyncFailedException` and + :class:`~singer_sdk.exceptions.AbortedSyncPausedException` are *not* + caught here — they abort the entire run and propagate to :meth:`invoke`. + + Raises: + AbortedSyncFailedException: If a lifecycle abort (non-resumable) occurs. + AbortedSyncPausedException: If a lifecycle pause (resumable) occurs. + """ self._reset_state_progress_markers() self._set_compatible_replication_methods() if self.state: self._state_writer.write_state(self.state) stream: Stream - for stream in self.streams.values(): - if not stream.selected and not stream.has_selected_descendents: - self.logger.info("Skipping deselected stream '%s'.", stream.name) - continue - - if stream.parent_stream_type: - self.logger.debug( - "Child stream '%s' is expected to be called " - "by parent stream '%s'. " - "Skipping direct invocation.", - type(stream).__name__, - stream.parent_stream_type.__name__, - ) - continue - - stream.sync() - stream.finalize_state_progress_markers() - - # this second loop is needed for all streams to print out their costs - # including child streams which are otherwise skipped in the loop above - for stream in self.streams.values(): - stream.log_sync_costs() + try: + for stream in self.streams.values(): + if not stream.selected and not stream.has_selected_descendents: + self.logger.info("Skipping deselected stream '%s'.", stream.name) + continue + + if stream.parent_stream_type: + self.logger.debug( + "Child stream '%s' is expected to be called " + "by parent stream '%s'. " + "Skipping direct invocation.", + type(stream).__name__, + stream.parent_stream_type.__name__, + ) + continue + + try: + stream.sync() + except AbortedSyncPausedException: + # Graceful pause: sync reached a valid resumable state. + # Finalize partial progress before propagating to invoke(). + stream.finalize_state_progress_markers() + raise + except AbortedSyncFailedException: + # Non-resumable abort: state is not stable — do NOT finalize. + raise + except Exception: + # stream.sync_result is already FAILED (set inside Stream.sync()). + # Log and continue to the next stream. + self.logger.exception( + "Stream '%s' failed; continuing with remaining streams.", + stream.name, + ) + else: + # Only reached when stream.sync() did not raise — SUCCESS. + stream.finalize_state_progress_markers() + finally: + # Always log results and costs — runs even when an abort exception + # propagates out of the per-stream loop. + for stream in self.streams.values(): + self._log_stream_sync_result(stream) + stream.log_sync_costs() # Command Line Execution @@ -552,7 +613,28 @@ def invoke( # type: ignore[override] parse_env_config=config.parse_env, validate_config=True, ) - tap.sync_all() + try: + tap.sync_all() + except AbortedSyncPausedException: + # State already written by _abort_sync(). Sync is resumable → exit 0. + sys.exit(0) + except AbortedSyncFailedException: + # Sync aborted; state is not stable → exit 1. + sys.exit(1) + + # Check for per-stream fatal errors (non-lifecycle). + failed_streams = [ + name + for name, stream in tap.streams.items() + if stream.sync_result is SyncResult.FAILED + ] + if failed_streams: + tap.logger.error( + "Sync completed with %d failed stream(s): %s", + len(failed_streams), + ", ".join(failed_streams), + ) + sys.exit(1) @classmethod def cb_discover( diff --git a/tests/core/snapshots/test_parent_child/test_child_deselected_parent/stderr.log b/tests/core/snapshots/test_parent_child/test_child_deselected_parent/stderr.log index beb7c98cb2..f8cdd5177a 100644 --- a/tests/core/snapshots/test_parent_child/test_child_deselected_parent/stderr.log +++ b/tests/core/snapshots/test_parent_child/test_child_deselected_parent/stderr.log @@ -2,3 +2,5 @@ INFO my-tap.parent Beginning sync of 'parent' in full_table mode INFO my-tap.child Beginning sync of 'child' in full_table mode with context: {'pid': 1} INFO my-tap.child Beginning sync of 'child' in full_table mode with context: {'pid': 2} INFO my-tap.child Beginning sync of 'child' in full_table mode with context: {'pid': 3} +INFO my-tap Stream 'child' sync result: success +INFO my-tap Stream 'parent' sync result: success diff --git a/tests/core/snapshots/test_parent_child/test_deselected_child/stderr.log b/tests/core/snapshots/test_parent_child/test_deselected_child/stderr.log index 026be91b9e..c778cd89a9 100644 --- a/tests/core/snapshots/test_parent_child/test_deselected_child/stderr.log +++ b/tests/core/snapshots/test_parent_child/test_deselected_child/stderr.log @@ -1,2 +1,3 @@ INFO my-tap Skipping deselected stream 'child'. INFO my-tap.parent Beginning sync of 'parent' in full_table mode +INFO my-tap Stream 'parent' sync result: success diff --git a/tests/core/snapshots/test_parent_child/test_one_parent_many_children/stderr.log b/tests/core/snapshots/test_parent_child/test_one_parent_many_children/stderr.log index 2abd97a272..625696351c 100644 --- a/tests/core/snapshots/test_parent_child/test_one_parent_many_children/stderr.log +++ b/tests/core/snapshots/test_parent_child/test_one_parent_many_children/stderr.log @@ -5,3 +5,5 @@ INFO my-tap-many.child_many Beginning sync of 'child_many' in full_table mode wi WARNING my-tap-many.child_many Properties ('composite_id', 'child_id') were present in the 'child_many' stream but not found in catalog schema. Ignoring. INFO my-tap-many.child_many Beginning sync of 'child_many' in full_table mode with context: {'child_id': 2, 'pid': '1'} INFO my-tap-many.child_many Beginning sync of 'child_many' in full_table mode with context: {'child_id': 3, 'pid': '1'} +INFO my-tap-many Stream 'child_many' sync result: success +INFO my-tap-many Stream 'parent_many' sync result: success diff --git a/tests/core/snapshots/test_parent_child/test_parent_context_fields_in_child/stderr.log b/tests/core/snapshots/test_parent_child/test_parent_context_fields_in_child/stderr.log index beb7c98cb2..f8cdd5177a 100644 --- a/tests/core/snapshots/test_parent_child/test_parent_context_fields_in_child/stderr.log +++ b/tests/core/snapshots/test_parent_child/test_parent_context_fields_in_child/stderr.log @@ -2,3 +2,5 @@ INFO my-tap.parent Beginning sync of 'parent' in full_table mode INFO my-tap.child Beginning sync of 'child' in full_table mode with context: {'pid': 1} INFO my-tap.child Beginning sync of 'child' in full_table mode with context: {'pid': 2} INFO my-tap.child Beginning sync of 'child' in full_table mode with context: {'pid': 3} +INFO my-tap Stream 'child' sync result: success +INFO my-tap Stream 'parent' sync result: success diff --git a/tests/core/snapshots/test_parent_child/test_preprocess_context_removes_large_payload/stderr.log b/tests/core/snapshots/test_parent_child/test_preprocess_context_removes_large_payload/stderr.log index 3949c11234..9e548df2a8 100644 --- a/tests/core/snapshots/test_parent_child/test_preprocess_context_removes_large_payload/stderr.log +++ b/tests/core/snapshots/test_parent_child/test_preprocess_context_removes_large_payload/stderr.log @@ -3,3 +3,5 @@ INFO tap-preprocess Added 'child_preprocessed' as child stream to 'parent_large' INFO tap-preprocess.parent_large Beginning sync of 'parent_large' in full_table mode INFO tap-preprocess.child_preprocessed Beginning sync of 'child_preprocessed' in full_table mode with context: {'parent_id': 1, 'parent_name': 'Parent A'} INFO tap-preprocess.child_preprocessed Beginning sync of 'child_preprocessed' in full_table mode with context: {'parent_id': 2, 'parent_name': 'Parent B'} +INFO tap-preprocess Stream 'child_preprocessed' sync result: success +INFO tap-preprocess Stream 'parent_large' sync result: success diff --git a/tests/core/test_sync_outcomes.py b/tests/core/test_sync_outcomes.py new file mode 100644 index 0000000000..ec3146f97c --- /dev/null +++ b/tests/core/test_sync_outcomes.py @@ -0,0 +1,214 @@ +"""Tests for per-stream sync outcome tracking and tap-level exit codes.""" + +from __future__ import annotations + +import json +import typing as t +from pathlib import Path + +import pytest +from click.testing import CliRunner + +from singer_sdk import Stream, Tap +from singer_sdk.exceptions import ( + AbortedSyncFailedException, + AbortedSyncPausedException, + FatalSyncError, +) +from singer_sdk.streams.core import SyncResult + +if t.TYPE_CHECKING: + from singer_sdk.helpers import types + + +# --------------------------------------------------------------------------- +# Stream fixtures +# --------------------------------------------------------------------------- + + +class GoodStream(Stream): + name = "good" + schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} + + def get_records(self, _context: types.Context | None): + yield {"id": 1} + + +class BadStream(Stream): + name = "bad" + schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} + + def get_records(self, _context: types.Context | None): + msg = "intentional failure" + raise FatalSyncError(msg) + + +class AbortPausedStream(Stream): + name = "abort_paused" + schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} + + def get_records(self, _context: types.Context | None): + raise AbortedSyncPausedException + + +class AbortFailedStream(Stream): + name = "abort_failed" + schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} + + def get_records(self, _context: types.Context | None): + msg = "forced" + raise AbortedSyncFailedException(msg) + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def make_tap(*stream_classes: type[Stream]) -> Tap: + """Create a minimal Tap instance with the given stream classes.""" + + class _Tap(Tap): + name = "test-outcomes-tap" + config_jsonschema: t.ClassVar = {"type": "object", "properties": {}} + + def discover_streams(self) -> list[Stream]: + return [cls(self) for cls in stream_classes] + + return _Tap(config={}) + + +def make_tap_class(*stream_classes: type[Stream]) -> type[Tap]: + """Return a Tap subclass (not an instance) for CLI invocation.""" + + class _Tap(Tap): + name = "test-outcomes-tap" + config_jsonschema: t.ClassVar = {"type": "object", "properties": {}} + + def discover_streams(self) -> list[Stream]: + return [cls(self) for cls in stream_classes] + + return _Tap + + +# --------------------------------------------------------------------------- +# Unit tests — sync_result attribute +# --------------------------------------------------------------------------- + + +def test_sync_result_default_is_none() -> None: + tap = make_tap(GoodStream) + stream = tap.streams["good"] + assert stream.sync_result is None + + +def test_sync_result_success() -> None: + tap = make_tap(GoodStream) + tap.sync_all() + assert tap.streams["good"].sync_result is SyncResult.SUCCESS + + +def test_sync_result_failed() -> None: + tap = make_tap(BadStream) + # sync_all() must NOT raise for per-stream fatal errors + tap.sync_all() + assert tap.streams["bad"].sync_result is SyncResult.FAILED + + +def test_sync_result_aborted_on_paused() -> None: + tap = make_tap(AbortPausedStream) + with pytest.raises(AbortedSyncPausedException): + tap.sync_all() + assert tap.streams["abort_paused"].sync_result is SyncResult.ABORTED + + +def test_sibling_continues_after_failure() -> None: + tap = make_tap(BadStream, GoodStream) + # Must not raise despite BadStream failing + tap.sync_all() + assert tap.streams["bad"].sync_result is SyncResult.FAILED + assert tap.streams["good"].sync_result is SyncResult.SUCCESS + + +def test_lifecycle_signal_stops_siblings() -> None: + tap = make_tap(AbortPausedStream, GoodStream) + with pytest.raises(AbortedSyncPausedException): + tap.sync_all() + # GoodStream never ran because the abort propagated before it was reached + assert tap.streams["good"].sync_result is None + + +# --------------------------------------------------------------------------- +# Unit tests — summary logging +# --------------------------------------------------------------------------- + + +def test_summary_logged_success(caplog: pytest.LogCaptureFixture) -> None: + tap = make_tap(GoodStream) + with caplog.at_level("INFO", logger="root"): + tap.sync_all() + assert "Stream 'good' sync result: success" in caplog.text + + +def test_summary_logged_failed(caplog: pytest.LogCaptureFixture) -> None: + tap = make_tap(BadStream) + with caplog.at_level("ERROR", logger="root"): + tap.sync_all() + assert "Stream 'bad' sync result: failed" in caplog.text + + +def test_summary_logged_aborted(caplog: pytest.LogCaptureFixture) -> None: + tap = make_tap(AbortPausedStream) + with ( + caplog.at_level("WARNING", logger="root"), + pytest.raises(AbortedSyncPausedException), + ): + tap.sync_all() + # Summary is emitted in the finally block even when sync_all raises + assert "Stream 'abort_paused' sync result: aborted" in caplog.text + + +def test_summary_not_logged_for_never_synced( + caplog: pytest.LogCaptureFixture, +) -> None: + """A stream that is skipped (deselected) must not produce a sync result line.""" + tap = make_tap(GoodStream) + stream = tap.streams["good"] + # Patch selected / has_selected_descendents so sync_all skips this stream + type(stream).selected = property(lambda _: False) # type: ignore[assignment] + type(stream).has_selected_descendents = property( # type: ignore[assignment] + lambda _: False + ) + with caplog.at_level("INFO", logger="root"): + tap.sync_all() + assert "sync result" not in caplog.text + + +# --------------------------------------------------------------------------- +# CLI exit code tests +# --------------------------------------------------------------------------- + + +def _cli_invoke(tap_cls: type[Tap]) -> int: + """Invoke tap CLI with an empty config file; return the exit code.""" + runner = CliRunner() + with runner.isolated_filesystem(): + Path("config.json").write_text(json.dumps({}), encoding="utf-8") + result = runner.invoke(tap_cls.cli, ["--config", "config.json"]) + return result.exit_code + + +def test_cli_exit_0_on_success() -> None: + assert _cli_invoke(make_tap_class(GoodStream)) == 0 + + +def test_cli_exit_1_on_stream_failure() -> None: + assert _cli_invoke(make_tap_class(BadStream)) == 1 + + +def test_cli_exit_0_on_aborted_sync_paused() -> None: + assert _cli_invoke(make_tap_class(AbortPausedStream)) == 0 + + +def test_cli_exit_1_on_aborted_sync_failed() -> None: + assert _cli_invoke(make_tap_class(AbortFailedStream)) == 1 diff --git a/tests/packages/snapshots/test_target_csv/test_countries_to_csv/singer.log b/tests/packages/snapshots/test_target_csv/test_countries_to_csv/singer.log index 32807b3c96..bc41a154f9 100644 --- a/tests/packages/snapshots/test_target_csv/test_countries_to_csv/singer.log +++ b/tests/packages/snapshots/test_target_csv/test_countries_to_csv/singer.log @@ -2,3 +2,5 @@ INFO tap-countries Skipping parse of env var settings... INFO target-csv Skipping parse of env var settings... INFO tap-countries.continents Beginning sync of 'continents' in full_table mode INFO tap-countries.countries Beginning sync of 'countries' in full_table mode +INFO tap-countries Stream 'continents' sync result: success +INFO tap-countries Stream 'countries' sync result: success diff --git a/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/activate_version/singer.log b/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/activate_version/singer.log index 3bfd61dceb..2875058ab3 100644 --- a/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/activate_version/singer.log +++ b/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/activate_version/singer.log @@ -4,6 +4,8 @@ INFO mapper-custom Skipping parse of env var settings... INFO mapper-custom Found '__else__=None' default mapper. Unmapped streams will be excluded from output. INFO tap-countries.continents Beginning sync of 'continents' in full_table mode INFO tap-countries.countries Beginning sync of 'countries' in full_table mode +INFO tap-countries Stream 'continents' sync result: success +INFO tap-countries Stream 'countries' sync result: success INFO mapper-custom Reader 'mapper-custom' completed processing 263 lines of input (2 schemas, 257 records, 0 batch manifests, 2 state messages, 2 activate version messages). WARNING target-csv The `ACTIVATE_VERSION` feature uses the `_sdc_deleted_at` and `_sdc_deleted_at` metadata properties so they will be added to the schema for '%s' even though `add_record_metadata` is disabled. WARNING target-csv.continents ACTIVATE_VERSION message received but not implemented by this target. Ignoring. diff --git a/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/no_activate_version/singer.log b/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/no_activate_version/singer.log index b8c05b936f..1a225257c3 100644 --- a/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/no_activate_version/singer.log +++ b/tests/packages/snapshots/test_target_csv/test_countries_to_csv_mapped/no_activate_version/singer.log @@ -3,5 +3,7 @@ INFO target-csv Skipping parse of env var settings... INFO mapper-custom Skipping parse of env var settings... INFO tap-countries.continents Beginning sync of 'continents' in full_table mode INFO tap-countries.countries Beginning sync of 'countries' in full_table mode +INFO tap-countries Stream 'continents' sync result: success +INFO tap-countries Stream 'countries' sync result: success INFO mapper-custom Reader 'mapper-custom' completed processing 261 lines of input (2 schemas, 257 records, 0 batch manifests, 2 state messages, 0 activate version messages). INFO target-csv Reader 'target-csv' completed processing 261 lines of input (2 schemas, 257 records, 0 batch manifests, 2 state messages, 0 activate version messages). diff --git a/tests/packages/snapshots/test_target_csv/test_fake_people_to_csv/singer.log b/tests/packages/snapshots/test_target_csv/test_fake_people_to_csv/singer.log index 5d920c8dce..b56bfab346 100644 --- a/tests/packages/snapshots/test_target_csv/test_fake_people_to_csv/singer.log +++ b/tests/packages/snapshots/test_target_csv/test_fake_people_to_csv/singer.log @@ -1,3 +1,4 @@ INFO tap-fake-people Skipping parse of env var settings... INFO target-csv Skipping parse of env var settings... INFO tap-fake-people.people Beginning sync of 'people' in full_table mode +INFO tap-fake-people Stream 'people' sync result: success From 5dc0393bfe4f748da3479693f0034b98fb0a4fc1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 12:12:04 -0600 Subject: [PATCH 02/17] refactor: Annotate `Context` parameter type as a mutable mapping MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/helpers/types.py | 4 ++-- singer_sdk/streams/core.py | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/singer_sdk/helpers/types.py b/singer_sdk/helpers/types.py index f01a21b9e6..e2a7690050 100644 --- a/singer_sdk/helpers/types.py +++ b/singer_sdk/helpers/types.py @@ -4,7 +4,7 @@ import os import typing as t -from collections.abc import Mapping +from collections.abc import MutableMapping import requests @@ -13,7 +13,7 @@ "Record", ] -Context: t.TypeAlias = Mapping[str, t.Any] +Context: t.TypeAlias = MutableMapping[str, t.Any] Record: t.TypeAlias = dict[str, t.Any] Auth: t.TypeAlias = t.Callable[[requests.PreparedRequest], requests.PreparedRequest] RequestFunc: t.TypeAlias = t.Callable[ diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index ebeee28ce4..10d2b6c79f 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -154,7 +154,7 @@ def __init__( self._logger: logging.Logger = tap.logger.getChild(self.name) self.metrics_logger = tap.metrics_logger self.tap_name: str = tap.name - self.context: types.Context | None = None + self.context: MappingProxyType | None = None self._config: dict = dict(tap.config) self._tap = tap From 4cbc37739ca1a973617fea70970cf5667d87253f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 12:13:04 -0600 Subject: [PATCH 03/17] chore: Update types in parent-child test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- pyproject.toml | 2 +- tests/core/test_parent_child.py | 48 ++++++++++++++++++++++----------- 2 files changed, 34 insertions(+), 16 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 616997ab46..6b0285ba2b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -48,7 +48,7 @@ dependencies = [ "simpleeval>=0.9.13,!=1.0.1", "simplejson>=3.17.6", "sqlalchemy>=2", - "typing-extensions>=4.5.0; python_version < '3.13'", + "typing-extensions>=4.5.0 ; python_full_version < '3.13'", "universal-pathlib>=0.2.6", "urllib3>=1.26.20", ] diff --git a/tests/core/test_parent_child.py b/tests/core/test_parent_child.py index d738a51b16..811c3dd065 100644 --- a/tests/core/test_parent_child.py +++ b/tests/core/test_parent_child.py @@ -3,6 +3,7 @@ import datetime import io import logging +import sys import typing as t from contextlib import redirect_stdout @@ -11,9 +12,16 @@ from singer_sdk import Stream, Tap +if sys.version_info >= (3, 12): + from typing import override # noqa: ICN003 +else: + from typing_extensions import override + if t.TYPE_CHECKING: from pytest_snapshot.plugin import Snapshot + from singer_sdk.helpers.types import Context, Record + DATETIME = datetime.datetime(2022, 1, 1, tzinfo=datetime.timezone.utc) @@ -28,17 +36,19 @@ class Parent(Stream): }, } + @override def get_child_context( self, - record: dict, - context: dict | None, # noqa: ARG002 - ) -> dict | None: + record: Record, + context: Context | None, + ) -> Context | None: """Create context for children streams.""" return { "pid": record["id"], } - def get_records(self, context: dict | None): # noqa: ARG002 + @override + def get_records(self, context: Context | None): """Get dummy records.""" yield {"id": 1} yield {"id": 2} @@ -58,7 +68,8 @@ class Child(Stream): } parent_stream_type = Parent - def get_records(self, context: dict | None): # noqa: ARG002 + @override + def get_records(self, context: Context | None): """Get dummy records.""" yield {"id": 1} yield {"id": 2} @@ -218,17 +229,19 @@ class ParentMany(Stream): }, } + @override def get_records( self, - context: dict | None, # noqa: ARG002 + context: Context | None, ) -> t.Iterable[dict | tuple[dict, dict | None]]: yield {"id": "1", "children": [1, 2, 3]} + @override def generate_child_contexts( self, - record: dict, - context: dict | None, # noqa: ARG002 - ) -> t.Iterable[dict | None]: + record: Record, + context: Context | None, + ) -> t.Iterable[Context | None]: for child_id in record["children"]: yield {"child_id": child_id, "pid": record["id"]} @@ -245,7 +258,8 @@ class ChildMany(Stream): } parent_stream_type = ParentMany - def get_records(self, context: dict | None): + @override + def get_records(self, context: Context | None): """Get dummy records.""" assert context is not None @@ -302,10 +316,11 @@ class ParentWithLargePayload(Stream): }, } + @override def get_child_context( self, - record: dict, - context: dict | None, # noqa: ARG002 + record: Record, + context: Context | None, ) -> dict | None: """Create context with large payload for child streams.""" return { @@ -315,7 +330,8 @@ def get_child_context( "large_payload": list(range(1, 1001)), } - def get_records(self, context: dict | None): # noqa: ARG002 + @override + def get_records(self, context: Context | None): """Get dummy records.""" yield {"id": 1, "name": "Parent A"} yield {"id": 2, "name": "Parent B"} @@ -342,12 +358,14 @@ def __init__(self, *args: t.Any, **kwargs: t.Any) -> None: def set_numbers(self, numbers: list[int]) -> None: self._numbers = numbers - def preprocess_context(self, context: dict) -> dict: + @override + def preprocess_context(self, context: Context) -> Context: """Remove large payload from parent context.""" self.set_numbers(context.pop("large_payload", [])) return context - def get_records(self, context: dict | None): + @override + def get_records(self, context: Context | None): """Get dummy records.""" # Verify that large_payload was removed assert context is not None From f80df500cdcf8a270679bb71b0746b2569f83a13 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 13:21:51 -0600 Subject: [PATCH 04/17] refactor: Reraise stream errors as sync abortions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/core.py | 36 +++++-- singer_sdk/tap_base.py | 42 +++----- tests/core/test_continue_on_errors.py | 148 ++++++++++++++++++++++++++ tests/core/test_sync_outcomes.py | 15 ++- 4 files changed, 195 insertions(+), 46 deletions(-) create mode 100644 tests/core/test_continue_on_errors.py diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index 10d2b6c79f..404521793d 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -1368,29 +1368,47 @@ def sync(self, context: types.Context | None = None) -> None: self._stream_version = self._initialized_at // 1000 self._write_activate_version_message(self._stream_version) + try: + self._run_sync(context) + except AbortedSyncFailedException: + self.sync_result = SyncResult.FAILED + raise + except AbortedSyncPausedException: + self.sync_result = SyncResult.ABORTED + raise + else: + self.sync_result = SyncResult.SUCCESS + + def _run_sync(self, context: types.Context | None) -> None: + """Execute the sync body, converting any non-lifecycle exception to one. + + Either completes normally or raises a lifecycle exception + (:class:`~singer_sdk.exceptions.AbortedSyncFailedException` or + :class:`~singer_sdk.exceptions.AbortedSyncPausedException`). + + Args: + context: Stream partition or context dictionary. + + Raises: + AbortedSyncFailedException: If the sync could not reach a resumable state. + AbortedSyncPausedException: If the sync paused at a resumable checkpoint. + """ try: batch_config = self.get_batch_config(self.config) if batch_config: self._sync_batches(batch_config, context=context) else: - # Sync the records themselves: for _ in self._sync_records(context=context): pass except (AbortedSyncFailedException, AbortedSyncPausedException): - # Lifecycle abort — must not be swallowed. Mark ABORTED and re-raise - # so sync_all() can propagate to invoke() for proper exit-code handling. - self.sync_result = SyncResult.ABORTED raise - except Exception: + except Exception as exc: # noqa: BLE001 self.log( "An unhandled error occurred while syncing '%s'", self.name, level=logging.ERROR, ) - self.sync_result = SyncResult.FAILED - raise - else: - self.sync_result = SyncResult.SUCCESS + self._abort_sync(exc) # always raises def _sync_children(self, child_context: types.Context | None) -> None: if child_context is None: diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index 5fd27e5753..37e965ce70 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -502,16 +502,16 @@ def _log_stream_sync_result(self, stream: Stream) -> None: def sync_all(self) -> None: """Sync all streams. - A stream that raises a fatal (non-lifecycle) exception is logged and - skipped; syncing continues with remaining streams. - - :class:`~singer_sdk.exceptions.AbortedSyncFailedException` and - :class:`~singer_sdk.exceptions.AbortedSyncPausedException` are *not* - caught here — they abort the entire run and propagate to :meth:`invoke`. - - Raises: - AbortedSyncFailedException: If a lifecycle abort (non-resumable) occurs. - AbortedSyncPausedException: If a lifecycle pause (resumable) occurs. + A stream that raises any exception is logged and skipped; syncing + continues with remaining streams. + + After all streams finish, streams whose :attr:`~singer_sdk.Stream.sync_result` + is :attr:`~singer_sdk.streams.core.SyncResult.FAILED` cause :meth:`invoke` + to exit with code 1. A + :class:`~singer_sdk.exceptions.AbortedSyncPausedException` is treated as + a graceful pause (exit 0); an + :class:`~singer_sdk.exceptions.AbortedSyncFailedException` is treated as + a failure (exit 1). """ self._reset_state_progress_markers() self._set_compatible_replication_methods() @@ -537,20 +537,13 @@ def sync_all(self) -> None: try: stream.sync() - except AbortedSyncPausedException: - # Graceful pause: sync reached a valid resumable state. - # Finalize partial progress before propagating to invoke(). - stream.finalize_state_progress_markers() - raise - except AbortedSyncFailedException: - # Non-resumable abort: state is not stable — do NOT finalize. - raise - except Exception: + except Exception as exc: # stream.sync_result is already FAILED (set inside Stream.sync()). # Log and continue to the next stream. - self.logger.exception( + self.logger.error( # noqa: TRY400 "Stream '%s' failed; continuing with remaining streams.", stream.name, + exc_info=exc.__cause__, ) else: # Only reached when stream.sync() did not raise — SUCCESS. @@ -613,14 +606,7 @@ def invoke( # type: ignore[override] parse_env_config=config.parse_env, validate_config=True, ) - try: - tap.sync_all() - except AbortedSyncPausedException: - # State already written by _abort_sync(). Sync is resumable → exit 0. - sys.exit(0) - except AbortedSyncFailedException: - # Sync aborted; state is not stable → exit 1. - sys.exit(1) + tap.sync_all() # Check for per-stream fatal errors (non-lifecycle). failed_streams = [ diff --git a/tests/core/test_continue_on_errors.py b/tests/core/test_continue_on_errors.py new file mode 100644 index 0000000000..d61a6d40e2 --- /dev/null +++ b/tests/core/test_continue_on_errors.py @@ -0,0 +1,148 @@ +from __future__ import annotations + +import sys +import typing as t +from datetime import datetime, timedelta, timezone + +from singer_sdk import Stream, Tap + +if sys.version_info >= (3, 12): + from typing import override # noqa: ICN003 +else: + from typing_extensions import override + + +if t.TYPE_CHECKING: + from singer_sdk.helpers.types import Context, Record + + +class _BaseStream(Stream): + max_records = 5 + fail_after = max_records + 1 + + schema: t.ClassVar = { + "type": "object", + "properties": { + "id": {"type": "integer"}, + "name": {"type": "string"}, + }, + } + + @override + def get_records(self, context: Context | None): + for i in range(1, self.max_records + 1): + yield {"id": i, "name": f"All Good ({i=})"} + + if i >= self.fail_after: + msg = "Something went wrong!" + raise RuntimeError(msg) + + +class _IncrementalBaseStream(_BaseStream): + replication_key = "updated_at" + schema: t.ClassVar = { + "type": "object", + "properties": { + "id": {"type": "integer"}, + "name": {"type": "string"}, + "updated_at": {"type": "string", "format": "date-time"}, + }, + } + + @override + def get_records(self, context: Context | None): + first_datetime = datetime(2024, 1, 1, tzinfo=timezone.utc) + start_date = self.get_starting_timestamp(context) + for i in range(1, self.max_records + 1): + rk = first_datetime + timedelta(minutes=i) + if start_date and rk <= start_date: + continue + + yield { + "id": i, + "name": f"All Good ({i=})", + "updated_at": rk.isoformat(), + } + + if i >= self.fail_after: + msg = "Something went wrong!" + raise RuntimeError(msg) + + +class StreamAllGood(_BaseStream): + """Stream that always succeeds.""" + + name = "all_good" + + +class StreamWithErrors(StreamAllGood): + """Stream that raises errors.""" + + name = "with_errors" + fail_after = 3 + + +class IncrementalAllGood(_IncrementalBaseStream): + """Incremental stream that raises errors.""" + + name = "incremental_all_good" + + +class IncrementalWithError(_IncrementalBaseStream): + """Incremental stream that raises errors.""" + + name = "incremental_with_errors" + fail_after = 3 + + +class IncrementalResumable(_IncrementalBaseStream): + """Incremental stream that raises errors.""" + + name = "incremental_resumable" + fail_after = 3 + is_sorted = True + + +class ParentStream(StreamAllGood): + """Parent stream that depends on the other streams.""" + + name = "parent" + + @override + def generate_child_contexts(self, record: Record, context: Context | None): + yield {"parent_id": record["id"]} + + +class ChildStreamWithErrors(StreamWithErrors): + """Child stream that raises errors.""" + + name = "child_with_errors" + schema: t.ClassVar = { + "type": "object", + "properties": { + "parent_id": {"type": "integer"}, + "id": {"type": "integer"}, + "name": {"type": "string"}, + }, + } + parent_stream_type = ParentStream + + +class ContinueOnErrorsTap(Tap): + name = "tap" + + @override + def discover_streams(self) -> list[Stream]: + return [ + StreamAllGood(self), + StreamWithErrors(self), + IncrementalAllGood(self), + IncrementalResumable(self), + IncrementalWithError(self), + ParentStream(self), + ChildStreamWithErrors(self), + ] + + +if __name__ == "__main__": + ContinueOnErrorsTap.cli() diff --git a/tests/core/test_sync_outcomes.py b/tests/core/test_sync_outcomes.py index ec3146f97c..0e87048878 100644 --- a/tests/core/test_sync_outcomes.py +++ b/tests/core/test_sync_outcomes.py @@ -6,7 +6,6 @@ import typing as t from pathlib import Path -import pytest from click.testing import CliRunner from singer_sdk import Stream, Tap @@ -18,6 +17,8 @@ from singer_sdk.streams.core import SyncResult if t.TYPE_CHECKING: + import pytest + from singer_sdk.helpers import types @@ -110,15 +111,13 @@ def test_sync_result_success() -> None: def test_sync_result_failed() -> None: tap = make_tap(BadStream) - # sync_all() must NOT raise for per-stream fatal errors tap.sync_all() assert tap.streams["bad"].sync_result is SyncResult.FAILED def test_sync_result_aborted_on_paused() -> None: tap = make_tap(AbortPausedStream) - with pytest.raises(AbortedSyncPausedException): - tap.sync_all() + tap.sync_all() assert tap.streams["abort_paused"].sync_result is SyncResult.ABORTED @@ -132,10 +131,9 @@ def test_sibling_continues_after_failure() -> None: def test_lifecycle_signal_stops_siblings() -> None: tap = make_tap(AbortPausedStream, GoodStream) - with pytest.raises(AbortedSyncPausedException): - tap.sync_all() - # GoodStream never ran because the abort propagated before it was reached - assert tap.streams["good"].sync_result is None + tap.sync_all() + # AbortPausedStream's exception is caught; GoodStream still runs successfully. + assert tap.streams["good"].sync_result is SyncResult.SUCCESS # --------------------------------------------------------------------------- @@ -161,7 +159,6 @@ def test_summary_logged_aborted(caplog: pytest.LogCaptureFixture) -> None: tap = make_tap(AbortPausedStream) with ( caplog.at_level("WARNING", logger="root"), - pytest.raises(AbortedSyncPausedException), ): tap.sync_all() # Summary is emitted in the finally block even when sync_all raises From addb3cc0b720c934262beb528bc5f81526216219 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 15:37:14 -0600 Subject: [PATCH 05/17] test: Add snapshot test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/core.py | 18 +++++- singer_sdk/tap_base.py | 9 ++- .../test_continue_on_errors/singer.jsonl | 57 +++++++++++++++++ .../singer_incrememtal.jsonl | 51 ++++++++++++++++ .../test_continue_on_errors/stderr.log | 31 ++++++++++ .../stderr_incremental.log | 28 +++++++++ tests/core/test_continue_on_errors.py | 61 ++++++++++++++++++- tests/core/test_sync_outcomes.py | 9 +-- 8 files changed, 255 insertions(+), 9 deletions(-) create mode 100644 tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer.jsonl create mode 100644 tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer_incrememtal.jsonl create mode 100644 tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log create mode 100644 tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index 404521793d..bbd9974d8d 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -1422,7 +1422,23 @@ def _sync_children(self, child_context: types.Context | None) -> None: for child_stream in self.child_streams: if child_stream.selected or child_stream.has_selected_descendents: - child_stream.sync(context=child_context) + try: + child_stream.sync(context=child_context) + except (AbortedSyncFailedException, AbortedSyncPausedException): + # sync_result already set inside child_stream.sync(). + # Continue syncing remaining children and let the parent + # record be written normally. + pass + except Exception as exc: # noqa: BLE001 + # Safety net — should not normally occur since _run_sync + # converts all non-lifecycle exceptions before they escape. + child_stream.sync_result = SyncResult.FAILED + self.log( + "Child stream '%s' failed unexpectedly.", + child_stream.name, + level=logging.ERROR, + exc_info=exc, + ) # Overridable Methods diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index 37e965ce70..dd8161ab8b 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -537,13 +537,16 @@ def sync_all(self) -> None: try: stream.sync() + except (AbortedSyncFailedException, AbortedSyncPausedException): + # sync_result is already set inside Stream.sync(). + # Summary is logged in the finally block; continue to next stream. + pass except Exception as exc: - # stream.sync_result is already FAILED (set inside Stream.sync()). - # Log and continue to the next stream. + # Unexpected non-lifecycle error (safety net). self.logger.error( # noqa: TRY400 "Stream '%s' failed; continuing with remaining streams.", stream.name, - exc_info=exc.__cause__, + exc_info=exc, ) else: # Only reached when stream.sync() did not raise — SUCCESS. diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer.jsonl b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer.jsonl new file mode 100644 index 0000000000..04c7e8f46b --- /dev/null +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer.jsonl @@ -0,0 +1,57 @@ +{"type":"SCHEMA","stream":"all_good","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"all_good","record":{"id":1,"name":"All Good (i=1)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":2,"name":"All Good (i=2)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":3,"name":"All Good (i=3)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":4,"name":"All Good (i=4)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":5,"name":"All Good (i=5)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"all_good":{}}}} +{"type":"SCHEMA","stream":"incremental_all_good","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"},"updated_at":{"format":"date-time","type":"string"}},"type":"object"},"key_properties":[],"bookmark_properties":["updated_at"]} +{"type":"RECORD","stream":"incremental_all_good","record":{"id":1,"name":"All Good (i=1)","updated_at":"2024-01-01T00:01:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_all_good","record":{"id":2,"name":"All Good (i=2)","updated_at":"2024-01-01T00:02:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_all_good","record":{"id":3,"name":"All Good (i=3)","updated_at":"2024-01-01T00:03:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_all_good","record":{"id":4,"name":"All Good (i=4)","updated_at":"2024-01-01T00:04:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_all_good","record":{"id":5,"name":"All Good (i=5)","updated_at":"2024-01-01T00:05:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"all_good":{},"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"}}}} +{"type":"SCHEMA","stream":"incremental_resumable","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"},"updated_at":{"format":"date-time","type":"string"}},"type":"object"},"key_properties":[],"bookmark_properties":["updated_at"]} +{"type":"RECORD","stream":"incremental_resumable","record":{"id":1,"name":"All Good (i=1)","updated_at":"2024-01-01T00:01:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_resumable","record":{"id":2,"name":"All Good (i=2)","updated_at":"2024-01-01T00:02:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_resumable","record":{"id":3,"name":"All Good (i=3)","updated_at":"2024-01-01T00:03:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"all_good":{},"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}}}} +{"type":"SCHEMA","stream":"incremental_with_errors","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"},"updated_at":{"format":"date-time","type":"string"}},"type":"object"},"key_properties":[],"bookmark_properties":["updated_at"]} +{"type":"RECORD","stream":"incremental_with_errors","record":{"id":1,"name":"All Good (i=1)","updated_at":"2024-01-01T00:01:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_with_errors","record":{"id":2,"name":"All Good (i=2)","updated_at":"2024-01-01T00:02:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_with_errors","record":{"id":3,"name":"All Good (i=3)","updated_at":"2024-01-01T00:03:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"all_good":{},"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}}}}} +{"type":"SCHEMA","stream":"parent","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":1},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":1},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":1},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"all_good":{},"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}},"child_with_errors":{}}}} +{"type":"RECORD","stream":"parent","record":{"id":1,"name":"All Good (i=1)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":2},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":2},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":2},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":2,"name":"All Good (i=2)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":3},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":3},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":3},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":3,"name":"All Good (i=3)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":4},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":4},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":4},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":4,"name":"All Good (i=4)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":5},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":5},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":5},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":5,"name":"All Good (i=5)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"all_good":{},"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}},"child_with_errors":{},"parent":{}}}} +{"type":"SCHEMA","stream":"with_errors","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"with_errors","record":{"id":1,"name":"All Good (i=1)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"with_errors","record":{"id":2,"name":"All Good (i=2)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"with_errors","record":{"id":3,"name":"All Good (i=3)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"all_good":{},"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}},"child_with_errors":{},"parent":{},"with_errors":{}}}} diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer_incrememtal.jsonl b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer_incrememtal.jsonl new file mode 100644 index 0000000000..8779923828 --- /dev/null +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/singer_incrememtal.jsonl @@ -0,0 +1,51 @@ +{"type":"STATE","value":{"bookmarks":{"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null}}}} +{"type":"SCHEMA","stream":"all_good","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"all_good","record":{"id":1,"name":"All Good (i=1)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":2,"name":"All Good (i=2)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":3,"name":"All Good (i=3)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":4,"name":"All Good (i=4)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"all_good","record":{"id":5,"name":"All Good (i=5)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null},"all_good":{}}}} +{"type":"SCHEMA","stream":"incremental_all_good","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"},"updated_at":{"format":"date-time","type":"string"}},"type":"object"},"key_properties":[],"bookmark_properties":["updated_at"]} +{"type":"SCHEMA","stream":"incremental_resumable","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"},"updated_at":{"format":"date-time","type":"string"}},"type":"object"},"key_properties":[],"bookmark_properties":["updated_at"]} +{"type":"RECORD","stream":"incremental_resumable","record":{"id":4,"name":"All Good (i=4)","updated_at":"2024-01-01T00:04:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_resumable","record":{"id":5,"name":"All Good (i=5)","updated_at":"2024-01-01T00:05:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null},"all_good":{}}}} +{"type":"SCHEMA","stream":"incremental_with_errors","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"},"updated_at":{"format":"date-time","type":"string"}},"type":"object"},"key_properties":[],"bookmark_properties":["updated_at"]} +{"type":"RECORD","stream":"incremental_with_errors","record":{"id":1,"name":"All Good (i=1)","updated_at":"2024-01-01T00:01:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_with_errors","record":{"id":2,"name":"All Good (i=2)","updated_at":"2024-01-01T00:02:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"incremental_with_errors","record":{"id":3,"name":"All Good (i=3)","updated_at":"2024-01-01T00:03:00+00:00"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}},"all_good":{}}}} +{"type":"SCHEMA","stream":"parent","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":1},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":1},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":1},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}},"all_good":{},"child_with_errors":{}}}} +{"type":"RECORD","stream":"parent","record":{"id":1,"name":"All Good (i=1)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":2},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":2},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":2},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":2,"name":"All Good (i=2)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":3},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":3},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":3},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":3,"name":"All Good (i=3)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":4},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":4},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":4},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":4,"name":"All Good (i=4)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"SCHEMA","stream":"child_with_errors","schema":{"properties":{"parent_id":{"type":"integer"},"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"child_with_errors","record":{"id":1,"name":"All Good (i=1)","parent_id":5},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":2,"name":"All Good (i=2)","parent_id":5},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"child_with_errors","record":{"id":3,"name":"All Good (i=3)","parent_id":5},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"parent","record":{"id":5,"name":"All Good (i=5)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}},"all_good":{},"child_with_errors":{},"parent":{}}}} +{"type":"SCHEMA","stream":"with_errors","schema":{"properties":{"id":{"type":"integer"},"name":{"type":"string"}},"type":"object"},"key_properties":[]} +{"type":"RECORD","stream":"with_errors","record":{"id":1,"name":"All Good (i=1)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"with_errors","record":{"id":2,"name":"All Good (i=2)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"RECORD","stream":"with_errors","record":{"id":3,"name":"All Good (i=3)"},"time_extracted":"2025-01-01T00:00:00+00:00"} +{"type":"STATE","value":{"bookmarks":{"incremental_all_good":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_resumable":{"replication_key":"updated_at","replication_key_value":"2024-01-01T00:05:00+00:00"},"incremental_with_errors":{"replication_key_signpost":"2025-01-01T00:00:00.000000+00:00","starting_replication_value":null,"progress_markers":{"Note":"Progress is not resumable if interrupted.","replication_key":"updated_at","replication_key_value":"2024-01-01T00:03:00+00:00"}},"all_good":{},"child_with_errors":{},"parent":{},"with_errors":{}}}} diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log new file mode 100644 index 0000000000..84ce9d1c44 --- /dev/null +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log @@ -0,0 +1,31 @@ +INFO tap Skipping parse of env var settings... +INFO tap Added 'child_with_errors' as child stream to 'parent' +INFO tap.all_good Beginning sync of 'all_good' in full_table mode +INFO tap.incremental_all_good Beginning sync of 'incremental_all_good' in incremental mode +INFO tap.incremental_all_good Starting incremental sync of 'incremental_all_good' with bookmark value: None +INFO tap.incremental_resumable Beginning sync of 'incremental_resumable' in incremental mode +INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumable' with bookmark value: None +ERROR tap.incremental_resumable An unhandled error occurred while syncing 'incremental_resumable' +INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode +INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None +ERROR tap.incremental_with_errors An unhandled error occurred while syncing 'incremental_with_errors' +INFO tap.parent Beginning sync of 'parent' in full_table mode +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 2} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 3} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 4} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 5} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode +ERROR tap.with_errors An unhandled error occurred while syncing 'with_errors' +INFO tap Stream 'all_good' sync result: success +ERROR tap Stream 'child_with_errors' sync result: failed +INFO tap Stream 'incremental_all_good' sync result: success +WARNING tap Stream 'incremental_resumable' sync result: aborted +ERROR tap Stream 'incremental_with_errors' sync result: failed +INFO tap Stream 'parent' sync result: success +ERROR tap Stream 'with_errors' sync result: failed diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log new file mode 100644 index 0000000000..b1db9eac14 --- /dev/null +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log @@ -0,0 +1,28 @@ +INFO tap.all_good Beginning sync of 'all_good' in full_table mode +INFO tap.incremental_all_good Beginning sync of 'incremental_all_good' in incremental mode +INFO tap.incremental_all_good Starting incremental sync of 'incremental_all_good' with bookmark value: 2024-01-01T00:05:00+00:00 +INFO tap.incremental_resumable Beginning sync of 'incremental_resumable' in incremental mode +INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumable' with bookmark value: 2024-01-01T00:03:00+00:00 +INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode +INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None +ERROR tap.incremental_with_errors An unhandled error occurred while syncing 'incremental_with_errors' +INFO tap.parent Beginning sync of 'parent' in full_table mode +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 2} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 3} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 4} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 5} +ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode +ERROR tap.with_errors An unhandled error occurred while syncing 'with_errors' +INFO tap Stream 'all_good' sync result: success +ERROR tap Stream 'child_with_errors' sync result: failed +INFO tap Stream 'incremental_all_good' sync result: success +INFO tap Stream 'incremental_resumable' sync result: success +ERROR tap Stream 'incremental_with_errors' sync result: failed +INFO tap Stream 'parent' sync result: success +ERROR tap Stream 'with_errors' sync result: failed diff --git a/tests/core/test_continue_on_errors.py b/tests/core/test_continue_on_errors.py index d61a6d40e2..ff17d0db8a 100644 --- a/tests/core/test_continue_on_errors.py +++ b/tests/core/test_continue_on_errors.py @@ -1,9 +1,16 @@ from __future__ import annotations +import io +import json +import logging import sys import typing as t +from contextlib import redirect_stdout from datetime import datetime, timedelta, timezone +import pytest +import time_machine + from singer_sdk import Stream, Tap if sys.version_info >= (3, 12): @@ -13,8 +20,12 @@ if t.TYPE_CHECKING: + from pytest_snapshot.plugin import Snapshot + from singer_sdk.helpers.types import Context, Record +DATETIME = datetime(2025, 1, 1, tzinfo=timezone.utc) + class _BaseStream(Stream): max_records = 5 @@ -53,6 +64,7 @@ class _IncrementalBaseStream(_BaseStream): def get_records(self, context: Context | None): first_datetime = datetime(2024, 1, 1, tzinfo=timezone.utc) start_date = self.get_starting_timestamp(context) + count = 0 for i in range(1, self.max_records + 1): rk = first_datetime + timedelta(minutes=i) if start_date and rk <= start_date: @@ -63,8 +75,9 @@ def get_records(self, context: Context | None): "name": f"All Good ({i=})", "updated_at": rk.isoformat(), } + count += 1 - if i >= self.fail_after: + if count >= self.fail_after: msg = "Something went wrong!" raise RuntimeError(msg) @@ -144,5 +157,51 @@ def discover_streams(self) -> list[Stream]: ] +@time_machine.travel(DATETIME, tick=False) +@pytest.mark.snapshot +def test_continue_on_errors( + caplog: pytest.LogCaptureFixture, + snapshot: Snapshot, +) -> None: + tap = ContinueOnErrorsTap() + buf = io.StringIO() + with ( + redirect_stdout(buf), + caplog.at_level("INFO"), + caplog.filtering(logging.Filter(tap.name)), + ): + tap.sync_all() + + buf.seek(0) + output = buf.read() + + snapshot.assert_match(output, "singer.jsonl") + snapshot.assert_match(caplog.text, "stderr.log") + + last_line = output.strip().splitlines()[-1] + state_message = json.loads(last_line) + state = state_message["value"] + bookmarks = state["bookmarks"] + + assert "replication_key_value" in bookmarks["incremental_resumable"] + assert "repllication_key_value" not in bookmarks["incremental_with_errors"] + + tap = ContinueOnErrorsTap(state=state) + buf = io.StringIO() + caplog.clear() + with ( + redirect_stdout(buf), + caplog.at_level("INFO"), + caplog.filtering(logging.Filter(tap.name)), + ): + tap.sync_all() + + buf.seek(0) + output = buf.read() + + snapshot.assert_match(output, "singer_incrememtal.jsonl") + snapshot.assert_match(caplog.text, "stderr_incremental.log") + + if __name__ == "__main__": ContinueOnErrorsTap.cli() diff --git a/tests/core/test_sync_outcomes.py b/tests/core/test_sync_outcomes.py index 0e87048878..acb6c8fdab 100644 --- a/tests/core/test_sync_outcomes.py +++ b/tests/core/test_sync_outcomes.py @@ -167,14 +167,15 @@ def test_summary_logged_aborted(caplog: pytest.LogCaptureFixture) -> None: def test_summary_not_logged_for_never_synced( caplog: pytest.LogCaptureFixture, + monkeypatch: pytest.MonkeyPatch, ) -> None: """A stream that is skipped (deselected) must not produce a sync result line.""" tap = make_tap(GoodStream) stream = tap.streams["good"] - # Patch selected / has_selected_descendents so sync_all skips this stream - type(stream).selected = property(lambda _: False) # type: ignore[assignment] - type(stream).has_selected_descendents = property( # type: ignore[assignment] - lambda _: False + stream_type = type(stream) + monkeypatch.setattr(stream_type, "selected", property(lambda _: False)) + monkeypatch.setattr( + stream_type, "has_selected_descendents", property(lambda _: False) ) with caplog.at_level("INFO", logger="root"): tap.sync_all() From c436f33cec3078b1093344c3a05c5387463f1ab1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 15:55:10 -0600 Subject: [PATCH 06/17] fix: Child failure means partial parent success MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/core.py | 43 +++++++++++++++++-- .../test_continue_on_errors/stderr.log | 2 +- .../stderr_incremental.log | 2 +- 3 files changed, 41 insertions(+), 6 deletions(-) diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index bbd9974d8d..af7164b666 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -76,6 +76,37 @@ class SyncResult(enum.Enum): PARTIAL = "partial" +def _combine_sync_results( + result1: SyncResult | None, + result2: SyncResult, +) -> SyncResult: + """Combine two SyncResults, treating None as SUCCESS. + + This is a helper function for combining a stream's existing sync_result with a new + result, treating None (the default state) as SUCCESS for combination purposes. + + Args: + result1: The first SyncResult, or None to treat as SUCCESS. + result2: The second SyncResult. + + Returns: + The combined SyncResult. + """ + if result1 is None: + return result2 + + if result1 is SyncResult.FAILED or result2 is SyncResult.FAILED: + return SyncResult.FAILED + + if result1 is SyncResult.ABORTED or result2 is SyncResult.ABORTED: + return SyncResult.ABORTED + + if result1 is SyncResult.PARTIAL or result2 is SyncResult.PARTIAL: + return SyncResult.PARTIAL + + return SyncResult.SUCCESS + + class Stream(abc.ABC): # noqa: PLR0904 """Abstract base class for tap streams. @@ -1377,7 +1408,11 @@ def sync(self, context: types.Context | None = None) -> None: self.sync_result = SyncResult.ABORTED raise else: - self.sync_result = SyncResult.SUCCESS + # Only mark SUCCESS if no child failure already degraded the result. + self.sync_result = _combine_sync_results( + self.sync_result, + SyncResult.SUCCESS, + ) def _run_sync(self, context: types.Context | None) -> None: """Execute the sync body, converting any non-lifecycle exception to one. @@ -1426,9 +1461,9 @@ def _sync_children(self, child_context: types.Context | None) -> None: child_stream.sync(context=child_context) except (AbortedSyncFailedException, AbortedSyncPausedException): # sync_result already set inside child_stream.sync(). - # Continue syncing remaining children and let the parent - # record be written normally. - pass + # Mark the parent failed too, then continue so remaining + # parent records and children are still attempted. + self.sync_result = SyncResult.PARTIAL except Exception as exc: # noqa: BLE001 # Safety net — should not normally occur since _run_sync # converts all non-lifecycle exceptions before they escape. diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log index 84ce9d1c44..55e0080d5c 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log @@ -27,5 +27,5 @@ ERROR tap Stream 'child_with_errors' sync result: failed INFO tap Stream 'incremental_all_good' sync result: success WARNING tap Stream 'incremental_resumable' sync result: aborted ERROR tap Stream 'incremental_with_errors' sync result: failed -INFO tap Stream 'parent' sync result: success +WARNING tap Stream 'parent' sync result: partial ERROR tap Stream 'with_errors' sync result: failed diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log index b1db9eac14..fd50fcc1fe 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log @@ -24,5 +24,5 @@ ERROR tap Stream 'child_with_errors' sync result: failed INFO tap Stream 'incremental_all_good' sync result: success INFO tap Stream 'incremental_resumable' sync result: success ERROR tap Stream 'incremental_with_errors' sync result: failed -INFO tap Stream 'parent' sync result: success +WARNING tap Stream 'parent' sync result: partial ERROR tap Stream 'with_errors' sync result: failed From 4de13842949f08006f25d4668da71be16e581855 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 16:15:30 -0600 Subject: [PATCH 07/17] refactor: Reorganize code MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/__init__.py | 3 +- singer_sdk/streams/_result.py | 86 ++++++++++++++++++ singer_sdk/streams/core.py | 56 +----------- singer_sdk/tap_base.py | 21 +---- tests/core/streams/test_sync_result.py | 119 +++++++++++++++++++++++++ 5 files changed, 211 insertions(+), 74 deletions(-) create mode 100644 singer_sdk/streams/_result.py create mode 100644 tests/core/streams/test_sync_result.py diff --git a/singer_sdk/streams/__init__.py b/singer_sdk/streams/__init__.py index a053233451..b82ca29a63 100644 --- a/singer_sdk/streams/__init__.py +++ b/singer_sdk/streams/__init__.py @@ -6,7 +6,8 @@ import warnings from singer_sdk.helpers._compat import SingerSDKDeprecationWarning -from singer_sdk.streams.core import Stream, SyncResult +from singer_sdk.streams._result import SyncResult +from singer_sdk.streams.core import Stream from singer_sdk.streams.graphql import GraphQLStream from singer_sdk.streams.rest import RESTStream diff --git a/singer_sdk/streams/_result.py b/singer_sdk/streams/_result.py new file mode 100644 index 0000000000..8b7757d08d --- /dev/null +++ b/singer_sdk/streams/_result.py @@ -0,0 +1,86 @@ +"""Sync outcome tracking for Singer streams.""" + +from __future__ import annotations + +import enum +import logging + +__all__ = ["SyncResult", "log_sync_result"] + + +class SyncResult(enum.Enum): + """Outcome of a single stream's sync operation. + + Set on :attr:`~singer_sdk.Stream.sync_result` after + :meth:`~singer_sdk.Stream.sync` completes or fails. + + Attributes: + SUCCESS: Completed without error. + FAILED: Raised a fatal (non-lifecycle) exception. + ABORTED: Raised a lifecycle abort exception + (:class:`~singer_sdk.exceptions.AbortedSyncFailedException` or + :class:`~singer_sdk.exceptions.AbortedSyncPausedException`). + PARTIAL: Reserved — ignorable errors with skipped records (requires PR 3). + """ + + SUCCESS = "success" + FAILED = "failed" + ABORTED = "aborted" + PARTIAL = "partial" + + @classmethod + def combine(cls, result1: SyncResult | None, result2: SyncResult) -> SyncResult: + """Merge two results; None means 'no prior result'. + + Priority order: FAILED > ABORTED > PARTIAL > SUCCESS. + + Args: + result1: The first SyncResult, or None to treat as no prior result. + result2: The second SyncResult. + + Returns: + The combined SyncResult. + """ + if result1 is None: + return result2 + + if result1 is cls.FAILED or result2 is cls.FAILED: + return cls.FAILED + + if result1 is cls.ABORTED or result2 is cls.ABORTED: + return cls.ABORTED + + if result1 is cls.PARTIAL or result2 is cls.PARTIAL: + return cls.PARTIAL + + return cls.SUCCESS + + +_LEVEL_MAP: dict[SyncResult, int] = { + SyncResult.SUCCESS: logging.INFO, + SyncResult.FAILED: logging.ERROR, + SyncResult.ABORTED: logging.WARNING, + SyncResult.PARTIAL: logging.WARNING, +} + + +def log_sync_result( + logger: logging.Logger, + stream_name: str, + result: SyncResult | None, +) -> None: + """Log a one-line sync outcome. No-op when result is None. + + Args: + logger: The logger to use. + stream_name: The name of the stream. + result: The sync result, or None if the stream was never synced. + """ + if result is None: + return + logger.log( + _LEVEL_MAP.get(result, logging.INFO), + "Stream '%s' sync result: %s", + stream_name, + result.value, + ) diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index af7164b666..a40a1d4a01 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -5,7 +5,6 @@ import abc import copy import datetime -import enum import json import logging import typing as t @@ -44,6 +43,7 @@ REPLICATION_INCREMENTAL, REPLICATION_LOG_BASED, # noqa: F401 ) +from singer_sdk.streams._result import SyncResult from singer_sdk.streams._state import StreamStateManager if t.TYPE_CHECKING: @@ -55,58 +55,6 @@ from singer_sdk.tap_base import Tap -class SyncResult(enum.Enum): - """Outcome of a single stream's sync operation. - - Set on :attr:`~singer_sdk.Stream.sync_result` after - :meth:`~singer_sdk.Stream.sync` completes or fails. - - Attributes: - SUCCESS: Completed without error. - FAILED: Raised a fatal (non-lifecycle) exception. - ABORTED: Raised a lifecycle abort exception - (:class:`~singer_sdk.exceptions.AbortedSyncFailedException` or - :class:`~singer_sdk.exceptions.AbortedSyncPausedException`). - PARTIAL: Reserved — ignorable errors with skipped records (requires PR 3). - """ - - SUCCESS = "success" - FAILED = "failed" - ABORTED = "aborted" - PARTIAL = "partial" - - -def _combine_sync_results( - result1: SyncResult | None, - result2: SyncResult, -) -> SyncResult: - """Combine two SyncResults, treating None as SUCCESS. - - This is a helper function for combining a stream's existing sync_result with a new - result, treating None (the default state) as SUCCESS for combination purposes. - - Args: - result1: The first SyncResult, or None to treat as SUCCESS. - result2: The second SyncResult. - - Returns: - The combined SyncResult. - """ - if result1 is None: - return result2 - - if result1 is SyncResult.FAILED or result2 is SyncResult.FAILED: - return SyncResult.FAILED - - if result1 is SyncResult.ABORTED or result2 is SyncResult.ABORTED: - return SyncResult.ABORTED - - if result1 is SyncResult.PARTIAL or result2 is SyncResult.PARTIAL: - return SyncResult.PARTIAL - - return SyncResult.SUCCESS - - class Stream(abc.ABC): # noqa: PLR0904 """Abstract base class for tap streams. @@ -1409,7 +1357,7 @@ def sync(self, context: types.Context | None = None) -> None: raise else: # Only mark SUCCESS if no child failure already degraded the result. - self.sync_result = _combine_sync_results( + self.sync_result = SyncResult.combine( self.sync_result, SyncResult.SUCCESS, ) diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index dd8161ab8b..164b66ad29 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -5,7 +5,6 @@ import abc import collections.abc import contextlib -import logging import sys import typing as t import warnings @@ -31,7 +30,7 @@ from singer_sdk.io_base import SingerWriter from singer_sdk.plugin_base import BaseSingerWriter, PluginBase, _ConfigInput from singer_sdk.singerlib import Catalog -from singer_sdk.streams.core import SyncResult +from singer_sdk.streams._result import SyncResult, log_sync_result if t.TYPE_CHECKING: from pathlib import PurePath @@ -480,23 +479,7 @@ def _log_stream_sync_result(self, stream: Stream) -> None: Args: stream: The stream whose result to log. """ - result = stream.sync_result - if result is None: - # Stream was never directly synced (deselected or child stream). - return - - level_map = { - SyncResult.SUCCESS: logging.INFO, - SyncResult.FAILED: logging.ERROR, - SyncResult.ABORTED: logging.WARNING, - SyncResult.PARTIAL: logging.WARNING, - } - self.logger.log( - level_map.get(result, logging.INFO), - "Stream '%s' sync result: %s", - stream.name, - result.value, - ) + log_sync_result(self.logger, stream.name, stream.sync_result) @t.final def sync_all(self) -> None: diff --git a/tests/core/streams/test_sync_result.py b/tests/core/streams/test_sync_result.py new file mode 100644 index 0000000000..0921800bda --- /dev/null +++ b/tests/core/streams/test_sync_result.py @@ -0,0 +1,119 @@ +"""Unit tests for singer_sdk.streams._result.""" + +from __future__ import annotations + +import logging +import unittest.mock + +import pytest + +from singer_sdk.streams._result import SyncResult, log_sync_result + + +class TestSyncResultCombine: + """Tests for SyncResult.combine.""" + + @pytest.mark.parametrize( + "result2", + [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], + ) + def test_none_returns_result2(self, result2: SyncResult) -> None: + assert SyncResult.combine(None, result2) is result2 + + @pytest.mark.parametrize( + "other", + [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], + ) + def test_failed_beats_all(self, other: SyncResult) -> None: + assert SyncResult.combine(SyncResult.FAILED, other) is SyncResult.FAILED + assert SyncResult.combine(other, SyncResult.FAILED) is SyncResult.FAILED + + @pytest.mark.parametrize( + "other", + [SyncResult.SUCCESS, SyncResult.PARTIAL], + ) + def test_aborted_beats_success_and_partial(self, other: SyncResult) -> None: + assert SyncResult.combine(SyncResult.ABORTED, other) is SyncResult.ABORTED + assert SyncResult.combine(other, SyncResult.ABORTED) is SyncResult.ABORTED + + def test_aborted_loses_to_failed(self) -> None: + assert ( + SyncResult.combine(SyncResult.ABORTED, SyncResult.FAILED) + is SyncResult.FAILED + ) + assert ( + SyncResult.combine(SyncResult.FAILED, SyncResult.ABORTED) + is SyncResult.FAILED + ) + + def test_partial_beats_success(self) -> None: + assert ( + SyncResult.combine(SyncResult.PARTIAL, SyncResult.SUCCESS) + is SyncResult.PARTIAL + ) + assert ( + SyncResult.combine(SyncResult.SUCCESS, SyncResult.PARTIAL) + is SyncResult.PARTIAL + ) + + def test_partial_loses_to_failed_and_aborted(self) -> None: + assert ( + SyncResult.combine(SyncResult.PARTIAL, SyncResult.FAILED) + is SyncResult.FAILED + ) + assert ( + SyncResult.combine(SyncResult.PARTIAL, SyncResult.ABORTED) + is SyncResult.ABORTED + ) + + @pytest.mark.parametrize( + "value", + [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], + ) + def test_same_value_is_idempotent(self, value: SyncResult) -> None: + assert SyncResult.combine(value, value) is value + + +class TestLogSyncResult: + """Tests for log_sync_result.""" + + def test_none_result_does_not_call_logger(self) -> None: + logger = unittest.mock.MagicMock(spec=logging.Logger) + log_sync_result(logger, "my_stream", None) + logger.log.assert_not_called() + + @pytest.mark.parametrize( + "result, expected_level", + [ + (SyncResult.SUCCESS, logging.INFO), + (SyncResult.FAILED, logging.ERROR), + (SyncResult.ABORTED, logging.WARNING), + (SyncResult.PARTIAL, logging.WARNING), + ], + ) + def test_correct_log_level( + self, + result: SyncResult, + expected_level: int, + ) -> None: + logger = unittest.mock.MagicMock(spec=logging.Logger) + log_sync_result(logger, "my_stream", result) + logger.log.assert_called_once() + call_args = logger.log.call_args + assert call_args[0][0] == expected_level + + @pytest.mark.parametrize( + "result", + [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], + ) + def test_message_contains_stream_name_and_result(self, result: SyncResult) -> None: + logger = unittest.mock.MagicMock(spec=logging.Logger) + log_sync_result(logger, "my_stream", result) + call_args = logger.log.call_args + # The format string and positional args are passed to logger.log + fmt = call_args[0][1] + name_arg = call_args[0][2] + value_arg = call_args[0][3] + assert name_arg == "my_stream" + assert result.value == value_arg + assert "%s" in fmt From 417cd737de9218d14601debd920c0d2453c2327f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 16:22:41 -0600 Subject: [PATCH 08/17] refactor: Remove superfluous exception catches MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/core.py | 10 ---------- singer_sdk/tap_base.py | 9 +++------ .../test_continue_on_errors/stderr.log | 3 +++ .../test_continue_on_errors/stderr_incremental.log | 2 ++ 4 files changed, 8 insertions(+), 16 deletions(-) diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index a40a1d4a01..e033815b13 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -1412,16 +1412,6 @@ def _sync_children(self, child_context: types.Context | None) -> None: # Mark the parent failed too, then continue so remaining # parent records and children are still attempted. self.sync_result = SyncResult.PARTIAL - except Exception as exc: # noqa: BLE001 - # Safety net — should not normally occur since _run_sync - # converts all non-lifecycle exceptions before they escape. - child_stream.sync_result = SyncResult.FAILED - self.log( - "Child stream '%s' failed unexpectedly.", - child_stream.name, - level=logging.ERROR, - exc_info=exc, - ) # Overridable Methods diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index 164b66ad29..13125363f9 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -520,16 +520,13 @@ def sync_all(self) -> None: try: stream.sync() - except (AbortedSyncFailedException, AbortedSyncPausedException): + except (AbortedSyncFailedException, AbortedSyncPausedException) as exc: # sync_result is already set inside Stream.sync(). # Summary is logged in the finally block; continue to next stream. - pass - except Exception as exc: - # Unexpected non-lifecycle error (safety net). self.logger.error( # noqa: TRY400 - "Stream '%s' failed; continuing with remaining streams.", + "Stream '%s' failed: %s", stream.name, - exc_info=exc, + exc.__cause__, ) else: # Only reached when stream.sync() did not raise — SUCCESS. diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log index 55e0080d5c..a4e45ea397 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log @@ -6,9 +6,11 @@ INFO tap.incremental_all_good Starting incremental sync of 'incremental_all_good INFO tap.incremental_resumable Beginning sync of 'incremental_resumable' in incremental mode INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumable' with bookmark value: None ERROR tap.incremental_resumable An unhandled error occurred while syncing 'incremental_resumable' +ERROR tap Stream 'incremental_resumable' failed: Something went wrong! INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None ERROR tap.incremental_with_errors An unhandled error occurred while syncing 'incremental_with_errors' +ERROR tap Stream 'incremental_with_errors' failed: Something went wrong! INFO tap.parent Beginning sync of 'parent' in full_table mode INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' @@ -22,6 +24,7 @@ INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table m ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode ERROR tap.with_errors An unhandled error occurred while syncing 'with_errors' +ERROR tap Stream 'with_errors' failed: Something went wrong! INFO tap Stream 'all_good' sync result: success ERROR tap Stream 'child_with_errors' sync result: failed INFO tap Stream 'incremental_all_good' sync result: success diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log index fd50fcc1fe..41f5a2c327 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log @@ -6,6 +6,7 @@ INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumab INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None ERROR tap.incremental_with_errors An unhandled error occurred while syncing 'incremental_with_errors' +ERROR tap Stream 'incremental_with_errors' failed: Something went wrong! INFO tap.parent Beginning sync of 'parent' in full_table mode INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' @@ -19,6 +20,7 @@ INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table m ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode ERROR tap.with_errors An unhandled error occurred while syncing 'with_errors' +ERROR tap Stream 'with_errors' failed: Something went wrong! INFO tap Stream 'all_good' sync result: success ERROR tap Stream 'child_with_errors' sync result: failed INFO tap Stream 'incremental_all_good' sync result: success From 58b31d98f2b6c3aba7c1aed93ae2d07efad5f56d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 16:23:57 -0600 Subject: [PATCH 09/17] Revert "refactor: Annotate `Context` parameter type as a mutable mapping" This reverts commit 5dc0393bfe4f748da3479693f0034b98fb0a4fc1. --- singer_sdk/helpers/types.py | 4 ++-- singer_sdk/streams/core.py | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/singer_sdk/helpers/types.py b/singer_sdk/helpers/types.py index e2a7690050..f01a21b9e6 100644 --- a/singer_sdk/helpers/types.py +++ b/singer_sdk/helpers/types.py @@ -4,7 +4,7 @@ import os import typing as t -from collections.abc import MutableMapping +from collections.abc import Mapping import requests @@ -13,7 +13,7 @@ "Record", ] -Context: t.TypeAlias = MutableMapping[str, t.Any] +Context: t.TypeAlias = Mapping[str, t.Any] Record: t.TypeAlias = dict[str, t.Any] Auth: t.TypeAlias = t.Callable[[requests.PreparedRequest], requests.PreparedRequest] RequestFunc: t.TypeAlias = t.Callable[ diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index e033815b13..4c5e8a835a 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -133,7 +133,7 @@ def __init__( self._logger: logging.Logger = tap.logger.getChild(self.name) self.metrics_logger = tap.metrics_logger self.tap_name: str = tap.name - self.context: MappingProxyType | None = None + self.context: types.Context | None = None self._config: dict = dict(tap.config) self._tap = tap From f23d3e10dc93324539f27cc78b031772a59bb550 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 16:32:13 -0600 Subject: [PATCH 10/17] refactor: Simplify a bit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/__init__.py | 3 +- singer_sdk/streams/_result.py | 2 - singer_sdk/tap_base.py | 79 +++++++++++++++------------------- 3 files changed, 35 insertions(+), 49 deletions(-) diff --git a/singer_sdk/streams/__init__.py b/singer_sdk/streams/__init__.py index b82ca29a63..f95dc12c2b 100644 --- a/singer_sdk/streams/__init__.py +++ b/singer_sdk/streams/__init__.py @@ -6,12 +6,11 @@ import warnings from singer_sdk.helpers._compat import SingerSDKDeprecationWarning -from singer_sdk.streams._result import SyncResult from singer_sdk.streams.core import Stream from singer_sdk.streams.graphql import GraphQLStream from singer_sdk.streams.rest import RESTStream -__all__ = ["GraphQLStream", "RESTStream", "Stream", "SyncResult"] +__all__ = ["GraphQLStream", "RESTStream", "Stream"] def __getattr__(name: str) -> t.Any: # noqa: ANN401 diff --git a/singer_sdk/streams/_result.py b/singer_sdk/streams/_result.py index 8b7757d08d..16e677905a 100644 --- a/singer_sdk/streams/_result.py +++ b/singer_sdk/streams/_result.py @@ -5,8 +5,6 @@ import enum import logging -__all__ = ["SyncResult", "log_sync_result"] - class SyncResult(enum.Enum): """Outcome of a single stream's sync operation. diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index 13125363f9..522c30ddf1 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -471,16 +471,6 @@ def _set_compatible_replication_methods(self) -> None: # Sync methods - def _log_stream_sync_result(self, stream: Stream) -> None: - """Log a one-line sync outcome for a stream. - - Called once per stream at the end of :meth:`sync_all`. - - Args: - stream: The stream whose result to log. - """ - log_sync_result(self.logger, stream.name, stream.sync_result) - @t.final def sync_all(self) -> None: """Sync all streams. @@ -502,41 +492,40 @@ def sync_all(self) -> None: self._state_writer.write_state(self.state) stream: Stream - try: - for stream in self.streams.values(): - if not stream.selected and not stream.has_selected_descendents: - self.logger.info("Skipping deselected stream '%s'.", stream.name) - continue - - if stream.parent_stream_type: - self.logger.debug( - "Child stream '%s' is expected to be called " - "by parent stream '%s'. " - "Skipping direct invocation.", - type(stream).__name__, - stream.parent_stream_type.__name__, - ) - continue - - try: - stream.sync() - except (AbortedSyncFailedException, AbortedSyncPausedException) as exc: - # sync_result is already set inside Stream.sync(). - # Summary is logged in the finally block; continue to next stream. - self.logger.error( # noqa: TRY400 - "Stream '%s' failed: %s", - stream.name, - exc.__cause__, - ) - else: - # Only reached when stream.sync() did not raise — SUCCESS. - stream.finalize_state_progress_markers() - finally: - # Always log results and costs — runs even when an abort exception - # propagates out of the per-stream loop. - for stream in self.streams.values(): - self._log_stream_sync_result(stream) - stream.log_sync_costs() + for stream in self.streams.values(): + if not stream.selected and not stream.has_selected_descendents: + self.logger.info("Skipping deselected stream '%s'.", stream.name) + continue + + if stream.parent_stream_type: + self.logger.debug( + "Child stream '%s' is expected to be called " + "by parent stream '%s'. " + "Skipping direct invocation.", + type(stream).__name__, + stream.parent_stream_type.__name__, + ) + continue + + try: + stream.sync() + except (AbortedSyncFailedException, AbortedSyncPausedException) as exc: + # sync_result is already set inside Stream.sync(). + # Summary is logged in the finally block; continue to next stream. + self.logger.error( # noqa: TRY400 + "Stream '%s' failed: %s", + stream.name, + exc.__cause__, + ) + else: + # Only reached when stream.sync() did not raise — SUCCESS. + stream.finalize_state_progress_markers() + + # Always log results and costs — runs even when an abort exception + # propagates out of the per-stream loop. + for stream in self.streams.values(): + log_sync_result(self.logger, stream.name, stream.sync_result) + stream.log_sync_costs() # Command Line Execution From 351ce2b8eb190351e8aef971b49ced53786ec625 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 16:40:34 -0600 Subject: [PATCH 11/17] test: Remove redundant tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- tests/core/test_sync_outcomes.py | 212 ------------------------------- 1 file changed, 212 deletions(-) delete mode 100644 tests/core/test_sync_outcomes.py diff --git a/tests/core/test_sync_outcomes.py b/tests/core/test_sync_outcomes.py deleted file mode 100644 index acb6c8fdab..0000000000 --- a/tests/core/test_sync_outcomes.py +++ /dev/null @@ -1,212 +0,0 @@ -"""Tests for per-stream sync outcome tracking and tap-level exit codes.""" - -from __future__ import annotations - -import json -import typing as t -from pathlib import Path - -from click.testing import CliRunner - -from singer_sdk import Stream, Tap -from singer_sdk.exceptions import ( - AbortedSyncFailedException, - AbortedSyncPausedException, - FatalSyncError, -) -from singer_sdk.streams.core import SyncResult - -if t.TYPE_CHECKING: - import pytest - - from singer_sdk.helpers import types - - -# --------------------------------------------------------------------------- -# Stream fixtures -# --------------------------------------------------------------------------- - - -class GoodStream(Stream): - name = "good" - schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} - - def get_records(self, _context: types.Context | None): - yield {"id": 1} - - -class BadStream(Stream): - name = "bad" - schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} - - def get_records(self, _context: types.Context | None): - msg = "intentional failure" - raise FatalSyncError(msg) - - -class AbortPausedStream(Stream): - name = "abort_paused" - schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} - - def get_records(self, _context: types.Context | None): - raise AbortedSyncPausedException - - -class AbortFailedStream(Stream): - name = "abort_failed" - schema: t.ClassVar = {"type": "object", "properties": {"id": {"type": "integer"}}} - - def get_records(self, _context: types.Context | None): - msg = "forced" - raise AbortedSyncFailedException(msg) - - -# --------------------------------------------------------------------------- -# Helpers -# --------------------------------------------------------------------------- - - -def make_tap(*stream_classes: type[Stream]) -> Tap: - """Create a minimal Tap instance with the given stream classes.""" - - class _Tap(Tap): - name = "test-outcomes-tap" - config_jsonschema: t.ClassVar = {"type": "object", "properties": {}} - - def discover_streams(self) -> list[Stream]: - return [cls(self) for cls in stream_classes] - - return _Tap(config={}) - - -def make_tap_class(*stream_classes: type[Stream]) -> type[Tap]: - """Return a Tap subclass (not an instance) for CLI invocation.""" - - class _Tap(Tap): - name = "test-outcomes-tap" - config_jsonschema: t.ClassVar = {"type": "object", "properties": {}} - - def discover_streams(self) -> list[Stream]: - return [cls(self) for cls in stream_classes] - - return _Tap - - -# --------------------------------------------------------------------------- -# Unit tests — sync_result attribute -# --------------------------------------------------------------------------- - - -def test_sync_result_default_is_none() -> None: - tap = make_tap(GoodStream) - stream = tap.streams["good"] - assert stream.sync_result is None - - -def test_sync_result_success() -> None: - tap = make_tap(GoodStream) - tap.sync_all() - assert tap.streams["good"].sync_result is SyncResult.SUCCESS - - -def test_sync_result_failed() -> None: - tap = make_tap(BadStream) - tap.sync_all() - assert tap.streams["bad"].sync_result is SyncResult.FAILED - - -def test_sync_result_aborted_on_paused() -> None: - tap = make_tap(AbortPausedStream) - tap.sync_all() - assert tap.streams["abort_paused"].sync_result is SyncResult.ABORTED - - -def test_sibling_continues_after_failure() -> None: - tap = make_tap(BadStream, GoodStream) - # Must not raise despite BadStream failing - tap.sync_all() - assert tap.streams["bad"].sync_result is SyncResult.FAILED - assert tap.streams["good"].sync_result is SyncResult.SUCCESS - - -def test_lifecycle_signal_stops_siblings() -> None: - tap = make_tap(AbortPausedStream, GoodStream) - tap.sync_all() - # AbortPausedStream's exception is caught; GoodStream still runs successfully. - assert tap.streams["good"].sync_result is SyncResult.SUCCESS - - -# --------------------------------------------------------------------------- -# Unit tests — summary logging -# --------------------------------------------------------------------------- - - -def test_summary_logged_success(caplog: pytest.LogCaptureFixture) -> None: - tap = make_tap(GoodStream) - with caplog.at_level("INFO", logger="root"): - tap.sync_all() - assert "Stream 'good' sync result: success" in caplog.text - - -def test_summary_logged_failed(caplog: pytest.LogCaptureFixture) -> None: - tap = make_tap(BadStream) - with caplog.at_level("ERROR", logger="root"): - tap.sync_all() - assert "Stream 'bad' sync result: failed" in caplog.text - - -def test_summary_logged_aborted(caplog: pytest.LogCaptureFixture) -> None: - tap = make_tap(AbortPausedStream) - with ( - caplog.at_level("WARNING", logger="root"), - ): - tap.sync_all() - # Summary is emitted in the finally block even when sync_all raises - assert "Stream 'abort_paused' sync result: aborted" in caplog.text - - -def test_summary_not_logged_for_never_synced( - caplog: pytest.LogCaptureFixture, - monkeypatch: pytest.MonkeyPatch, -) -> None: - """A stream that is skipped (deselected) must not produce a sync result line.""" - tap = make_tap(GoodStream) - stream = tap.streams["good"] - stream_type = type(stream) - monkeypatch.setattr(stream_type, "selected", property(lambda _: False)) - monkeypatch.setattr( - stream_type, "has_selected_descendents", property(lambda _: False) - ) - with caplog.at_level("INFO", logger="root"): - tap.sync_all() - assert "sync result" not in caplog.text - - -# --------------------------------------------------------------------------- -# CLI exit code tests -# --------------------------------------------------------------------------- - - -def _cli_invoke(tap_cls: type[Tap]) -> int: - """Invoke tap CLI with an empty config file; return the exit code.""" - runner = CliRunner() - with runner.isolated_filesystem(): - Path("config.json").write_text(json.dumps({}), encoding="utf-8") - result = runner.invoke(tap_cls.cli, ["--config", "config.json"]) - return result.exit_code - - -def test_cli_exit_0_on_success() -> None: - assert _cli_invoke(make_tap_class(GoodStream)) == 0 - - -def test_cli_exit_1_on_stream_failure() -> None: - assert _cli_invoke(make_tap_class(BadStream)) == 1 - - -def test_cli_exit_0_on_aborted_sync_paused() -> None: - assert _cli_invoke(make_tap_class(AbortPausedStream)) == 0 - - -def test_cli_exit_1_on_aborted_sync_failed() -> None: - assert _cli_invoke(make_tap_class(AbortFailedStream)) == 1 From b68335486fd051a5c1401beea56d1eba20a442d0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 16:51:03 -0600 Subject: [PATCH 12/17] refactor: Remove redundant log MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/tap_base.py | 17 +---------------- 1 file changed, 1 insertion(+), 16 deletions(-) diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index 522c30ddf1..d56f10e6c1 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -5,7 +5,6 @@ import abc import collections.abc import contextlib -import sys import typing as t import warnings from enum import Enum @@ -30,7 +29,7 @@ from singer_sdk.io_base import SingerWriter from singer_sdk.plugin_base import BaseSingerWriter, PluginBase, _ConfigInput from singer_sdk.singerlib import Catalog -from singer_sdk.streams._result import SyncResult, log_sync_result +from singer_sdk.streams._result import log_sync_result if t.TYPE_CHECKING: from pathlib import PurePath @@ -580,20 +579,6 @@ def invoke( # type: ignore[override] ) tap.sync_all() - # Check for per-stream fatal errors (non-lifecycle). - failed_streams = [ - name - for name, stream in tap.streams.items() - if stream.sync_result is SyncResult.FAILED - ] - if failed_streams: - tap.logger.error( - "Sync completed with %d failed stream(s): %s", - len(failed_streams), - ", ".join(failed_streams), - ) - sys.exit(1) - @classmethod def cb_discover( cls: type[Tap], From 90c8ce94228708796f9537bf3beb08edc871fbb5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 17:09:11 -0600 Subject: [PATCH 13/17] test: Simplify tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- tests/core/streams/test_sync_result.py | 157 ++++++++++--------------- 1 file changed, 62 insertions(+), 95 deletions(-) diff --git a/tests/core/streams/test_sync_result.py b/tests/core/streams/test_sync_result.py index 0921800bda..09bb72fa4c 100644 --- a/tests/core/streams/test_sync_result.py +++ b/tests/core/streams/test_sync_result.py @@ -3,117 +3,84 @@ from __future__ import annotations import logging -import unittest.mock import pytest from singer_sdk.streams._result import SyncResult, log_sync_result -class TestSyncResultCombine: - """Tests for SyncResult.combine.""" - - @pytest.mark.parametrize( - "result2", - [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], - ) - def test_none_returns_result2(self, result2: SyncResult) -> None: - assert SyncResult.combine(None, result2) is result2 - - @pytest.mark.parametrize( - "other", - [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], - ) - def test_failed_beats_all(self, other: SyncResult) -> None: - assert SyncResult.combine(SyncResult.FAILED, other) is SyncResult.FAILED - assert SyncResult.combine(other, SyncResult.FAILED) is SyncResult.FAILED - - @pytest.mark.parametrize( - "other", - [SyncResult.SUCCESS, SyncResult.PARTIAL], - ) - def test_aborted_beats_success_and_partial(self, other: SyncResult) -> None: - assert SyncResult.combine(SyncResult.ABORTED, other) is SyncResult.ABORTED - assert SyncResult.combine(other, SyncResult.ABORTED) is SyncResult.ABORTED - - def test_aborted_loses_to_failed(self) -> None: - assert ( - SyncResult.combine(SyncResult.ABORTED, SyncResult.FAILED) - is SyncResult.FAILED - ) - assert ( - SyncResult.combine(SyncResult.FAILED, SyncResult.ABORTED) - is SyncResult.FAILED +def test_sync_result_combine(subtests: pytest.Subtests) -> None: + """Test SyncResult.combine with all combinations of SyncResult values.""" + result_none = None + with subtests.test("None returns the other result"): + assert all( + SyncResult.combine(result_none, other_result) is other_result + for other_result in SyncResult ) - def test_partial_beats_success(self) -> None: - assert ( - SyncResult.combine(SyncResult.PARTIAL, SyncResult.SUCCESS) - is SyncResult.PARTIAL + with subtests.test("FAILED beats all"): + result = SyncResult.FAILED + assert all( + SyncResult.combine(result, other) is SyncResult.FAILED + for other in SyncResult ) - assert ( - SyncResult.combine(SyncResult.SUCCESS, SyncResult.PARTIAL) - is SyncResult.PARTIAL + assert all( + SyncResult.combine(other, result) is SyncResult.FAILED + for other in SyncResult ) - def test_partial_loses_to_failed_and_aborted(self) -> None: - assert ( - SyncResult.combine(SyncResult.PARTIAL, SyncResult.FAILED) - is SyncResult.FAILED + with subtests.test("ABORTED beats SUCCESS and PARTIAL"): + result = SyncResult.ABORTED + assert all( + SyncResult.combine(result, other) is SyncResult.ABORTED + for other in [SyncResult.SUCCESS, SyncResult.PARTIAL] ) - assert ( - SyncResult.combine(SyncResult.PARTIAL, SyncResult.ABORTED) - is SyncResult.ABORTED + assert all( + SyncResult.combine(other, result) is SyncResult.ABORTED + for other in [SyncResult.SUCCESS, SyncResult.PARTIAL] ) - @pytest.mark.parametrize( - "value", - [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], - ) - def test_same_value_is_idempotent(self, value: SyncResult) -> None: - assert SyncResult.combine(value, value) is value + with subtests.test("PARTIAL beats SUCCESS"): + result = SyncResult.PARTIAL + assert SyncResult.combine(result, SyncResult.SUCCESS) is SyncResult.PARTIAL + assert SyncResult.combine(SyncResult.SUCCESS, result) is SyncResult.PARTIAL + with subtests.test("Idempotent combinations"): + assert all( + SyncResult.combine(any_result, any_result) is any_result + for any_result in SyncResult + ) -class TestLogSyncResult: - """Tests for log_sync_result.""" - def test_none_result_does_not_call_logger(self) -> None: - logger = unittest.mock.MagicMock(spec=logging.Logger) +def test_none_result_does_not_call_logger(caplog: pytest.LogCaptureFixture) -> None: + logger = logging.getLogger("test_logger") + with caplog.at_level("INFO", logger="test_logger"): log_sync_result(logger, "my_stream", None) - logger.log.assert_not_called() - - @pytest.mark.parametrize( - "result, expected_level", - [ - (SyncResult.SUCCESS, logging.INFO), - (SyncResult.FAILED, logging.ERROR), - (SyncResult.ABORTED, logging.WARNING), - (SyncResult.PARTIAL, logging.WARNING), - ], - ) - def test_correct_log_level( - self, - result: SyncResult, - expected_level: int, - ) -> None: - logger = unittest.mock.MagicMock(spec=logging.Logger) - log_sync_result(logger, "my_stream", result) - logger.log.assert_called_once() - call_args = logger.log.call_args - assert call_args[0][0] == expected_level - - @pytest.mark.parametrize( - "result", - [SyncResult.SUCCESS, SyncResult.FAILED, SyncResult.ABORTED, SyncResult.PARTIAL], - ) - def test_message_contains_stream_name_and_result(self, result: SyncResult) -> None: - logger = unittest.mock.MagicMock(spec=logging.Logger) + + assert len(caplog.records) == 0 + + +@pytest.mark.parametrize( + "result, expected_level", + [ + (SyncResult.SUCCESS, logging.INFO), + (SyncResult.FAILED, logging.ERROR), + (SyncResult.ABORTED, logging.WARNING), + (SyncResult.PARTIAL, logging.WARNING), + ], +) +def test_log_result( + result: SyncResult, + caplog: pytest.LogCaptureFixture, + expected_level: int, +) -> None: + logger = logging.getLogger("test_logger") + with caplog.at_level("INFO", logger="test_logger"): log_sync_result(logger, "my_stream", result) - call_args = logger.log.call_args - # The format string and positional args are passed to logger.log - fmt = call_args[0][1] - name_arg = call_args[0][2] - value_arg = call_args[0][3] - assert name_arg == "my_stream" - assert result.value == value_arg - assert "%s" in fmt + + assert len(caplog.records) == 1 + + record = caplog.records[0] + assert record.name == "test_logger" + assert record.levelno == expected_level + assert record.getMessage() == f"Stream 'my_stream' sync result: {result.value}" From 7651893ec7c488116661f6a3aeb9446056be5c9e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 17:16:16 -0600 Subject: [PATCH 14/17] refactor: Make `combine` an instance method MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/_result.py | 24 ++++++++-------- singer_sdk/streams/core.py | 5 +--- tests/core/streams/test_sync_result.py | 38 +++++++++----------------- 3 files changed, 25 insertions(+), 42 deletions(-) diff --git a/singer_sdk/streams/_result.py b/singer_sdk/streams/_result.py index 16e677905a..53a3dc1aff 100644 --- a/singer_sdk/streams/_result.py +++ b/singer_sdk/streams/_result.py @@ -26,32 +26,30 @@ class SyncResult(enum.Enum): ABORTED = "aborted" PARTIAL = "partial" - @classmethod - def combine(cls, result1: SyncResult | None, result2: SyncResult) -> SyncResult: + def combine(self, other: SyncResult | None) -> SyncResult: """Merge two results; None means 'no prior result'. Priority order: FAILED > ABORTED > PARTIAL > SUCCESS. Args: - result1: The first SyncResult, or None to treat as no prior result. - result2: The second SyncResult. + other: The first SyncResult, or None to treat as no prior result. Returns: The combined SyncResult. """ - if result1 is None: - return result2 + if other is None: + return self - if result1 is cls.FAILED or result2 is cls.FAILED: - return cls.FAILED + if self is SyncResult.FAILED or other is SyncResult.FAILED: + return SyncResult.FAILED - if result1 is cls.ABORTED or result2 is cls.ABORTED: - return cls.ABORTED + if self is SyncResult.ABORTED or other is SyncResult.ABORTED: + return SyncResult.ABORTED - if result1 is cls.PARTIAL or result2 is cls.PARTIAL: - return cls.PARTIAL + if self is SyncResult.PARTIAL or other is SyncResult.PARTIAL: + return SyncResult.PARTIAL - return cls.SUCCESS + return SyncResult.SUCCESS _LEVEL_MAP: dict[SyncResult, int] = { diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index 4c5e8a835a..ea6c5c57d3 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -1357,10 +1357,7 @@ def sync(self, context: types.Context | None = None) -> None: raise else: # Only mark SUCCESS if no child failure already degraded the result. - self.sync_result = SyncResult.combine( - self.sync_result, - SyncResult.SUCCESS, - ) + self.sync_result = SyncResult.SUCCESS.combine(self.sync_result) def _run_sync(self, context: types.Context | None) -> None: """Execute the sync body, converting any non-lifecycle exception to one. diff --git a/tests/core/streams/test_sync_result.py b/tests/core/streams/test_sync_result.py index 09bb72fa4c..d4cf781cd0 100644 --- a/tests/core/streams/test_sync_result.py +++ b/tests/core/streams/test_sync_result.py @@ -13,43 +13,31 @@ def test_sync_result_combine(subtests: pytest.Subtests) -> None: """Test SyncResult.combine with all combinations of SyncResult values.""" result_none = None with subtests.test("None returns the other result"): - assert all( - SyncResult.combine(result_none, other_result) is other_result - for other_result in SyncResult - ) + assert all(result.combine(result_none) is result for result in SyncResult) with subtests.test("FAILED beats all"): - result = SyncResult.FAILED - assert all( - SyncResult.combine(result, other) is SyncResult.FAILED - for other in SyncResult - ) - assert all( - SyncResult.combine(other, result) is SyncResult.FAILED - for other in SyncResult - ) + failed = SyncResult.FAILED + assert all(result.combine(failed) is SyncResult.FAILED for result in SyncResult) + assert all(failed.combine(result) is SyncResult.FAILED for result in SyncResult) with subtests.test("ABORTED beats SUCCESS and PARTIAL"): - result = SyncResult.ABORTED + aborted = SyncResult.ABORTED assert all( - SyncResult.combine(result, other) is SyncResult.ABORTED - for other in [SyncResult.SUCCESS, SyncResult.PARTIAL] + result.combine(aborted) is SyncResult.ABORTED + for result in [SyncResult.SUCCESS, SyncResult.PARTIAL] ) assert all( - SyncResult.combine(other, result) is SyncResult.ABORTED - for other in [SyncResult.SUCCESS, SyncResult.PARTIAL] + aborted.combine(result) is SyncResult.ABORTED + for result in [SyncResult.SUCCESS, SyncResult.PARTIAL] ) with subtests.test("PARTIAL beats SUCCESS"): - result = SyncResult.PARTIAL - assert SyncResult.combine(result, SyncResult.SUCCESS) is SyncResult.PARTIAL - assert SyncResult.combine(SyncResult.SUCCESS, result) is SyncResult.PARTIAL + partial = SyncResult.PARTIAL + assert SyncResult.SUCCESS.combine(partial) is SyncResult.PARTIAL + assert partial.combine(SyncResult.SUCCESS) is SyncResult.PARTIAL with subtests.test("Idempotent combinations"): - assert all( - SyncResult.combine(any_result, any_result) is any_result - for any_result in SyncResult - ) + assert all(result.combine(result) is result for result in SyncResult) def test_none_result_does_not_call_logger(caplog: pytest.LogCaptureFixture) -> None: From f4aeaaf33adcd7ae3c30c8cd8f743dc2ef12db38 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 17:28:37 -0600 Subject: [PATCH 15/17] refactor: Mention error MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/_result.py | 8 +++++--- singer_sdk/streams/core.py | 3 ++- singer_sdk/tap_base.py | 5 ++--- .../test_continue_on_errors/stderr.log | 16 ++++++++-------- .../stderr_incremental.log | 14 +++++++------- 5 files changed, 24 insertions(+), 22 deletions(-) diff --git a/singer_sdk/streams/_result.py b/singer_sdk/streams/_result.py index 53a3dc1aff..4199c27b96 100644 --- a/singer_sdk/streams/_result.py +++ b/singer_sdk/streams/_result.py @@ -27,12 +27,14 @@ class SyncResult(enum.Enum): PARTIAL = "partial" def combine(self, other: SyncResult | None) -> SyncResult: - """Merge two results; None means 'no prior result'. + """Merge this result with a prior accumulated result. - Priority order: FAILED > ABORTED > PARTIAL > SUCCESS. + ``self`` is the new outcome being applied; ``other`` is whatever has + accumulated so far (``None`` means no prior result, so ``self`` wins + unconditionally). Priority order: FAILED > ABORTED > PARTIAL > SUCCESS. Args: - other: The first SyncResult, or None to treat as no prior result. + other: The accumulated result so far, or None if this is the first. Returns: The combined SyncResult. diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index ea6c5c57d3..83c91009da 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -1384,8 +1384,9 @@ def _run_sync(self, context: types.Context | None) -> None: raise except Exception as exc: # noqa: BLE001 self.log( - "An unhandled error occurred while syncing '%s'", + "An error occurred while syncing '%s': %s", self.name, + str(exc), level=logging.ERROR, ) self._abort_sync(exc) # always raises diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index d56f10e6c1..80ff4164f8 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -477,8 +477,7 @@ def sync_all(self) -> None: A stream that raises any exception is logged and skipped; syncing continues with remaining streams. - After all streams finish, streams whose :attr:`~singer_sdk.Stream.sync_result` - is :attr:`~singer_sdk.streams.core.SyncResult.FAILED` cause :meth:`invoke` + After all streams finish, streams that failed cause :meth:`invoke` to exit with code 1. A :class:`~singer_sdk.exceptions.AbortedSyncPausedException` is treated as a graceful pause (exit 0); an @@ -510,7 +509,7 @@ def sync_all(self) -> None: stream.sync() except (AbortedSyncFailedException, AbortedSyncPausedException) as exc: # sync_result is already set inside Stream.sync(). - # Summary is logged in the finally block; continue to next stream. + # Result is logged below after the loop; continue to next stream. self.logger.error( # noqa: TRY400 "Stream '%s' failed: %s", stream.name, diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log index a4e45ea397..9daa777690 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log @@ -5,25 +5,25 @@ INFO tap.incremental_all_good Beginning sync of 'incremental_all_good' in increm INFO tap.incremental_all_good Starting incremental sync of 'incremental_all_good' with bookmark value: None INFO tap.incremental_resumable Beginning sync of 'incremental_resumable' in incremental mode INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumable' with bookmark value: None -ERROR tap.incremental_resumable An unhandled error occurred while syncing 'incremental_resumable' +ERROR tap.incremental_resumable An error occurred while syncing 'incremental_resumable': Something went wrong! ERROR tap Stream 'incremental_resumable' failed: Something went wrong! INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None -ERROR tap.incremental_with_errors An unhandled error occurred while syncing 'incremental_with_errors' +ERROR tap.incremental_with_errors An error occurred while syncing 'incremental_with_errors': Something went wrong! ERROR tap Stream 'incremental_with_errors' failed: Something went wrong! INFO tap.parent Beginning sync of 'parent' in full_table mode INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 2} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 3} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 4} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 5} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode -ERROR tap.with_errors An unhandled error occurred while syncing 'with_errors' +ERROR tap.with_errors An error occurred while syncing 'with_errors': Something went wrong! ERROR tap Stream 'with_errors' failed: Something went wrong! INFO tap Stream 'all_good' sync result: success ERROR tap Stream 'child_with_errors' sync result: failed diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log index 41f5a2c327..d17b0f6ecb 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log @@ -5,21 +5,21 @@ INFO tap.incremental_resumable Beginning sync of 'incremental_resumable' in incr INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumable' with bookmark value: 2024-01-01T00:03:00+00:00 INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None -ERROR tap.incremental_with_errors An unhandled error occurred while syncing 'incremental_with_errors' +ERROR tap.incremental_with_errors An error occurred while syncing 'incremental_with_errors': Something went wrong! ERROR tap Stream 'incremental_with_errors' failed: Something went wrong! INFO tap.parent Beginning sync of 'parent' in full_table mode INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 2} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 3} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 4} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 5} -ERROR tap.child_with_errors An unhandled error occurred while syncing 'child_with_errors' +ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode -ERROR tap.with_errors An unhandled error occurred while syncing 'with_errors' +ERROR tap.with_errors An error occurred while syncing 'with_errors': Something went wrong! ERROR tap Stream 'with_errors' failed: Something went wrong! INFO tap Stream 'all_good' sync result: success ERROR tap Stream 'child_with_errors' sync result: failed From c5240e67e093d05f3179791727712dc2a8833c10 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez=20Mondrag=C3=B3n?= Date: Fri, 6 Mar 2026 17:46:18 -0600 Subject: [PATCH 16/17] fix: Handle exit codes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez Mondragón --- singer_sdk/streams/_result.py | 10 ++++++++++ singer_sdk/tap_base.py | 22 +++++++++++++--------- tests/core/test_continue_on_errors.py | 4 ---- 3 files changed, 23 insertions(+), 13 deletions(-) diff --git a/singer_sdk/streams/_result.py b/singer_sdk/streams/_result.py index 4199c27b96..85aa1dda3f 100644 --- a/singer_sdk/streams/_result.py +++ b/singer_sdk/streams/_result.py @@ -5,6 +5,8 @@ import enum import logging +logger = logging.getLogger("singer_sdk") + class SyncResult(enum.Enum): """Outcome of a single stream's sync operation. @@ -53,6 +55,14 @@ def combine(self, other: SyncResult | None) -> SyncResult: return SyncResult.SUCCESS + def exit_code(self) -> int: + """Return the appropriate exit code for this result.""" + match self: + case SyncResult.SUCCESS: + return 0 + case SyncResult.FAILED | SyncResult.PARTIAL | SyncResult.ABORTED: + return 1 + _LEVEL_MAP: dict[SyncResult, int] = { SyncResult.SUCCESS: logging.INFO, diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index 80ff4164f8..727520bdb0 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -5,6 +5,7 @@ import abc import collections.abc import contextlib +import sys import typing as t import warnings from enum import Enum @@ -29,7 +30,7 @@ from singer_sdk.io_base import SingerWriter from singer_sdk.plugin_base import BaseSingerWriter, PluginBase, _ConfigInput from singer_sdk.singerlib import Catalog -from singer_sdk.streams._result import log_sync_result +from singer_sdk.streams._result import SyncResult, log_sync_result if t.TYPE_CHECKING: from pathlib import PurePath @@ -471,19 +472,17 @@ def _set_compatible_replication_methods(self) -> None: # Sync methods @t.final - def sync_all(self) -> None: + def sync_all(self) -> SyncResult: """Sync all streams. A stream that raises any exception is logged and skipped; syncing continues with remaining streams. - After all streams finish, streams that failed cause :meth:`invoke` - to exit with code 1. A - :class:`~singer_sdk.exceptions.AbortedSyncPausedException` is treated as - a graceful pause (exit 0); an - :class:`~singer_sdk.exceptions.AbortedSyncFailedException` is treated as - a failure (exit 1). + Returns: + The combined SyncResult of all streams. """ + result = SyncResult.SUCCESS + self._reset_state_progress_markers() self._set_compatible_replication_methods() if self.state: @@ -519,12 +518,16 @@ def sync_all(self) -> None: # Only reached when stream.sync() did not raise — SUCCESS. stream.finalize_state_progress_markers() + result = result.combine(stream.sync_result) + # Always log results and costs — runs even when an abort exception # propagates out of the per-stream loop. for stream in self.streams.values(): log_sync_result(self.logger, stream.name, stream.sync_result) stream.log_sync_costs() + return result + # Command Line Execution def _handle_termination( # pragma: no cover @@ -576,7 +579,8 @@ def invoke( # type: ignore[override] parse_env_config=config.parse_env, validate_config=True, ) - tap.sync_all() + result = tap.sync_all() + sys.exit(result.exit_code()) @classmethod def cb_discover( diff --git a/tests/core/test_continue_on_errors.py b/tests/core/test_continue_on_errors.py index ff17d0db8a..c408ce105d 100644 --- a/tests/core/test_continue_on_errors.py +++ b/tests/core/test_continue_on_errors.py @@ -201,7 +201,3 @@ def test_continue_on_errors( snapshot.assert_match(output, "singer_incrememtal.jsonl") snapshot.assert_match(caplog.text, "stderr_incremental.log") - - -if __name__ == "__main__": - ContinueOnErrorsTap.cli() From 7836109d8800c956c58bffe9d480da2eebc47e31 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edgar=20Ram=C3=ADrez-Mondrag=C3=B3n?= Date: Fri, 13 Mar 2026 20:54:09 -0600 Subject: [PATCH 17/17] refactor: Return `SyncResult` MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Edgar Ramírez-Mondragón --- singer_sdk/streams/core.py | 20 ++++++------------- singer_sdk/tap_base.py | 17 +++------------- .../test_continue_on_errors/stderr.log | 3 --- .../stderr_incremental.log | 2 -- 4 files changed, 9 insertions(+), 33 deletions(-) diff --git a/singer_sdk/streams/core.py b/singer_sdk/streams/core.py index 9ffd91f921..c94bb4942e 100644 --- a/singer_sdk/streams/core.py +++ b/singer_sdk/streams/core.py @@ -1306,7 +1306,7 @@ def _sync_batches( # Public methods ("final", not recommended to be overridden) @t.final - def sync(self, context: types.Context | None = None) -> None: + def sync(self, context: types.Context | None = None) -> SyncResult: """Sync this stream. This method is internal to the SDK and should not need to be overridden. @@ -1314,9 +1314,8 @@ def sync(self, context: types.Context | None = None) -> None: Args: context: Stream partition or context dictionary. - Raises: - AbortedSyncFailedException: If the sync was aborted non-resumably. - AbortedSyncPausedException: If the sync was paused at a resumable point. + Returns: + The SyncResult for this stream. """ # Preprocess context before it's frozen context = self.preprocess_context(context) if context else None @@ -1351,13 +1350,12 @@ def sync(self, context: types.Context | None = None) -> None: self._run_sync(context) except AbortedSyncFailedException: self.sync_result = SyncResult.FAILED - raise except AbortedSyncPausedException: self.sync_result = SyncResult.ABORTED - raise else: # Only mark SUCCESS if no child failure already degraded the result. self.sync_result = SyncResult.SUCCESS.combine(self.sync_result) + return self.sync_result def _run_sync(self, context: types.Context | None) -> None: """Execute the sync body, converting any non-lifecycle exception to one. @@ -1403,15 +1401,9 @@ def _sync_children(self, child_context: types.Context | None) -> None: for child_stream in self.child_streams: if child_stream.selected or child_stream.has_selected_descendents: - try: - child_stream.sync(context=child_context) - except (AbortedSyncFailedException, AbortedSyncPausedException): - # Child stream was interrupted, continue with remaining children - # sync_result already set inside child_stream.sync(). - # Mark the parent failed too, then continue so remaining - # parent records and children are still attempted. + child_result = child_stream.sync(context=child_context) + if child_result != SyncResult.SUCCESS: self.sync_result = SyncResult.PARTIAL - continue # Overridable Methods diff --git a/singer_sdk/tap_base.py b/singer_sdk/tap_base.py index bd9c736f72..397bfa4a81 100644 --- a/singer_sdk/tap_base.py +++ b/singer_sdk/tap_base.py @@ -505,21 +505,10 @@ def sync_all(self) -> SyncResult: ) continue - try: - stream.sync() - except (AbortedSyncFailedException, AbortedSyncPausedException) as exc: - # sync_result is already set inside Stream.sync(). - # Result is logged below after the loop; continue to next stream. - self.logger.error( # noqa: TRY400 - "Stream '%s' failed: %s", - stream.name, - exc.__cause__, - ) - else: - # Only reached when stream.sync() did not raise — SUCCESS. + stream_result = stream.sync() + if stream_result == SyncResult.SUCCESS: stream.finalize_state_progress_markers() - - result = result.combine(stream.sync_result) + result = result.combine(stream_result) # Always log results and costs — runs even when an abort exception # propagates out of the per-stream loop. diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log index 9daa777690..b181f10235 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr.log @@ -6,11 +6,9 @@ INFO tap.incremental_all_good Starting incremental sync of 'incremental_all_good INFO tap.incremental_resumable Beginning sync of 'incremental_resumable' in incremental mode INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumable' with bookmark value: None ERROR tap.incremental_resumable An error occurred while syncing 'incremental_resumable': Something went wrong! -ERROR tap Stream 'incremental_resumable' failed: Something went wrong! INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None ERROR tap.incremental_with_errors An error occurred while syncing 'incremental_with_errors': Something went wrong! -ERROR tap Stream 'incremental_with_errors' failed: Something went wrong! INFO tap.parent Beginning sync of 'parent' in full_table mode INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! @@ -24,7 +22,6 @@ INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table m ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode ERROR tap.with_errors An error occurred while syncing 'with_errors': Something went wrong! -ERROR tap Stream 'with_errors' failed: Something went wrong! INFO tap Stream 'all_good' sync result: success ERROR tap Stream 'child_with_errors' sync result: failed INFO tap Stream 'incremental_all_good' sync result: success diff --git a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log index d17b0f6ecb..604492178f 100644 --- a/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log +++ b/tests/core/snapshots/test_continue_on_errors/test_continue_on_errors/stderr_incremental.log @@ -6,7 +6,6 @@ INFO tap.incremental_resumable Starting incremental sync of 'incremental_resumab INFO tap.incremental_with_errors Beginning sync of 'incremental_with_errors' in incremental mode INFO tap.incremental_with_errors Starting incremental sync of 'incremental_with_errors' with bookmark value: None ERROR tap.incremental_with_errors An error occurred while syncing 'incremental_with_errors': Something went wrong! -ERROR tap Stream 'incremental_with_errors' failed: Something went wrong! INFO tap.parent Beginning sync of 'parent' in full_table mode INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table mode with context: {'parent_id': 1} ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! @@ -20,7 +19,6 @@ INFO tap.child_with_errors Beginning sync of 'child_with_errors' in full_table m ERROR tap.child_with_errors An error occurred while syncing 'child_with_errors': Something went wrong! INFO tap.with_errors Beginning sync of 'with_errors' in full_table mode ERROR tap.with_errors An error occurred while syncing 'with_errors': Something went wrong! -ERROR tap Stream 'with_errors' failed: Something went wrong! INFO tap Stream 'all_good' sync result: success ERROR tap Stream 'child_with_errors' sync result: failed INFO tap Stream 'incremental_all_good' sync result: success