Repository navigation
feat: Attempt to sync all streams instead of crashing on the first error - #3545
Conversation
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Reviewer's GuideIntroduce SyncResult-based per-stream sync outcome tracking, adjust tap-level sync_all and CLI exit codes to continue syncing remaining streams while distinguishing lifecycle aborts, and update typing, tests, and dependency constraints accordingly. Sequence diagram for tap sync_all with per-stream SyncResult and exit codessequenceDiagram
actor User
participant CLI as tap_entrypoint
participant Tap as TapBase
participant Stream1 as Stream_orders
participant Stream2 as Stream_customers
User->>CLI: invoke()
CLI->>Tap: invoke()
Tap->>Tap: tap = TapClass()
Tap->>Tap: result = tap.sync_all()
rect rgb(235, 245, 255)
Tap->>Stream1: sync()
Stream1->>Stream1: _run_sync(context=None)
Note over Stream1: Non-lifecycle error occurs
Stream1-->>Stream1: _abort_sync(exc)
Stream1-->>Tap: raise AbortedSyncFailedException
Tap->>Tap: stream.sync_result already set = FAILED
Tap->>Tap: result = result.combine(FAILED)
end
rect rgb(235, 245, 255)
Tap->>Stream2: sync()
Stream2->>Stream2: _run_sync(context=None)
Stream2-->>Stream2: completes without error
Stream2->>Stream2: sync_result = SUCCESS.combine(None)
Stream2-->>Tap: return
Tap->>Tap: stream.finalize_state_progress_markers()
Tap->>Tap: result = result.combine(SUCCESS)
end
Tap->>Tap: for each stream: log_sync_result(logger, name, sync_result)
Tap->>Tap: stream.log_sync_costs()
Tap-->>CLI: return result (FAILED)
CLI->>CLI: code = result.exit_code() # 1
CLI-->>User: process exit with code 1
Updated class diagram for Stream, TapBase, and SyncResultclassDiagram
class SyncResult {
<<enum>>
SUCCESS
FAILED
ABORTED
PARTIAL
+SyncResult combine(other SyncResult)
+int exit_code()
}
class Stream {
-dict _schema
-dict~str,int~ _sync_costs
+list~Stream~ child_streams
+SyncResult sync_result
+sync(context Context)
-_run_sync(context Context)
-_sync_children(child_context Context)
+finalize_state_progress_markers()
+log_sync_costs()
}
class TapBase {
-dict streams
+SyncResult sync_all()
+invoke()
}
class log_sync_result {
<<function>>
+log_sync_result(logger Logger, stream_name str, result SyncResult)
}
TapBase "1" --> "*" Stream : manages
Stream --> SyncResult : uses
TapBase --> SyncResult : returns
log_sync_result --> SyncResult : logs
TapBase ..> log_sync_result : calls
File-Level Changes
Tips and commandsInteracting with Sourcery
Customizing Your ExperienceAccess your dashboard to:
Getting Help
|
Documentation build overview
Show files changed (10 files in total): 📝 10 modified | ➕ 0 added | ➖ 0 deleted
|
There was a problem hiding this comment.
Hey - I've found 2 issues, and left some high level feedback:
- The
Contexttype alias was changed toMutableMapping[str, Any]whileStream.contextis now aMappingProxyType | None; this makes the stored context effectively read‑only and inconsistent with the alias, so consider aligning these types (e.g., keepContextas a mapping and introduce a separate alias for mutable contexts where needed). - In
sync_all, the error log for failed streams usesexc_info=exc.__cause__, which will often beNone; if the intention is to capture the full traceback for the failure, you likely wantexc_info=excinstead.
Prompt for AI Agents
Please address the comments from this code review:
## Overall Comments
- The `Context` type alias was changed to `MutableMapping[str, Any]` while `Stream.context` is now a `MappingProxyType | None`; this makes the stored context effectively read‑only and inconsistent with the alias, so consider aligning these types (e.g., keep `Context` as a mapping and introduce a separate alias for mutable contexts where needed).
- In `sync_all`, the error log for failed streams uses `exc_info=exc.__cause__`, which will often be `None`; if the intention is to capture the full traceback for the failure, you likely want `exc_info=exc` instead.
## Individual Comments
### Comment 1
<location path="singer_sdk/tap_base.py" line_range="518-496" />
<code_context>
+ )
+ continue
+
+ try:
+ stream.sync()
+ except Exception as exc:
+ # stream.sync_result is already FAILED (set inside Stream.sync()).
</code_context>
<issue_to_address>
**issue (bug_risk):** Lifecycle abort exceptions from `stream.sync()` are being caught here despite the docstring stating they should propagate.
`Stream.sync()` now re-raises `AbortedSyncFailedException` / `AbortedSyncPausedException` so they can propagate, and the `sync_all` docstring says these should *not* be caught here. With `except Exception as exc:`, they are still treated as non-fatal and only logged, which breaks the documented lifecycle semantics. Please either re-raise these two exceptions explicitly in this block or narrow the `except` so they bypass the generic handler (e.g., a dedicated `except` for the abort exceptions before `Exception`).
</issue_to_address>
### Comment 2
<location path="tests/core/test_sync_outcomes.py" line_range="171-180" />
<code_context>
+ 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
+
+
</code_context>
<issue_to_address>
**suggestion (testing):** Avoid permanently patching type attributes in the deselected-stream test
In `test_summary_not_logged_for_never_synced`, `selected` and `has_selected_descendents` are reassigned on `type(stream)`, permanently mutating `GoodStream` for all later tests. This risks subtle cross-test interference if other tests rely on the default selection behavior. Instead, either use `monkeypatch.setattr` so the properties are restored automatically, or create a short-lived subclass of `GoodStream` with the overridden properties and pass that into `make_tap`.
Suggested implementation:
```python
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)
# Temporarily patch selected / has_selected_descendents so sync_all skips this stream
monkeypatch.setattr(
stream_type,
"selected",
property(lambda _: False),
raising=True,
)
monkeypatch.setattr(
stream_type,
"has_selected_descendents",
property(lambda _: False),
raising=True,
)
with caplog.at_level("INFO", logger="root"):
tap.sync_all()
assert "sync result" not in caplog.text
```
If `pytest.MonkeyPatch` is not already imported in this file (often it's just used via the fixture name without type annotations), you may need to add or adjust type hints, or remove the explicit `pytest.MonkeyPatch` annotation depending on your existing typing conventions for fixtures.
</issue_to_address>Help me be more useful! Please click 👍 or 👎 on each comment and I'll use the feedback to improve your reviews.
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## feat/safely-ignore-errors #3545 +/- ##
=============================================================
+ Coverage 93.74% 94.12% +0.37%
=============================================================
Files 73 74 +1
Lines 5897 6617 +720
Branches 724 893 +169
=============================================================
+ Hits 5528 6228 +700
- Misses 274 292 +18
- Partials 95 97 +2
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
6cedb45 to
f80df50
Compare
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
…ing" This reverts commit 5dc0393.
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
4b52595 to
f23d3e1
Compare
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
c717d12 to
c5240e6
Compare
There was a problem hiding this comment.
Hey - I've found 1 issue, and left some high level feedback:
- In
test_continue_on_errors, the assertionassert "repllication_key_value" not in bookmarks["incremental_with_errors"]has a typo in the key name and will not correctly guard against an incorrectly persistedreplication_key_valuebookmark. - In
Tap.sync_all, the error log for failed streams usesexc.__cause__, which will often beNone; loggingexc(orstr(exc)) instead would give a more reliable and informative message. - Several new snapshot filenames/labels (e.g.
singer_incrememtal.jsonl) appear to be misspelled; if unintentional, correcting these now will avoid confusion for future readers and tooling.
Prompt for AI Agents
Please address the comments from this code review:
## Overall Comments
- In `test_continue_on_errors`, the assertion `assert "repllication_key_value" not in bookmarks["incremental_with_errors"]` has a typo in the key name and will not correctly guard against an incorrectly persisted `replication_key_value` bookmark.
- In `Tap.sync_all`, the error log for failed streams uses `exc.__cause__`, which will often be `None`; logging `exc` (or `str(exc)`) instead would give a more reliable and informative message.
- Several new snapshot filenames/labels (e.g. `singer_incrememtal.jsonl`) appear to be misspelled; if unintentional, correcting these now will avoid confusion for future readers and tooling.
## Individual Comments
### Comment 1
<location path="singer_sdk/streams/core.py" line_range="1383-1392" />
<code_context>
for _ in self._sync_records(context=context):
pass
- except Exception:
+ except (AbortedSyncFailedException, AbortedSyncPausedException):
+ raise
+ except Exception as exc: # noqa: BLE001
</code_context>
<issue_to_address>
**suggestion (bug_risk):** Overwriting `self.sync_result` with `PARTIAL` may accidentally downgrade a more severe result.
If `_sync_children` runs after a more severe outcome has already been set (e.g. `FAILED` or `ABORTED`), assigning `SyncResult.PARTIAL` directly would overwrite that and incorrectly lower the overall severity. To avoid this, use the existing aggregation logic, for example `self.sync_result = SyncResult.PARTIAL.combine(self.sync_result)`, so more severe results are preserved.
Suggested implementation:
```python
self.sync_result = SyncResult.PARTIAL.combine(self.sync_result)
```
Because I only see a small portion of the file, you should:
1. Search `singer_sdk/streams/core.py` for any other occurrences of `self.sync_result = SyncResult.PARTIAL` and update them in the same way, using `SyncResult.PARTIAL.combine(self.sync_result)`.
2. Confirm that `SyncResult.combine` indeed keeps the highest severity when combining (the suggested order assumes it does). If its contract is different (e.g., argument precedence is reversed), adjust the call order accordingly, for example `self.sync_result = self.sync_result.combine(SyncResult.PARTIAL)`.
</issue_to_address>Help me be more useful! Please click 👍 or 👎 on each comment and I'll use the feedback to improve your reviews.
…ests (#3556) https://github.com/meltano/sdk/blob/4a5dd894d69cc068d140acea6e15fa4040831611/singer_sdk/tap_base.py#L284-L286 ## Related - #3029 - #3030 - #3037 - #3545 ## Summary by Sourcery Adjust tap dry-run behavior to cap record counts on all streams while ensuring parent streams continue emitting records when child streams hit their record limits. Bug Fixes: - Prevent auto-generated dry-run syncs from aborting parent streams when child streams reach their record limit by catching child abort exceptions and stopping only remaining child syncs. - Apply the dry-run record limit consistently to both parent and child streams so all streams are subject to the same cap. ## Summary by Sourcery Ensure dry-run syncs apply record limits consistently across parent and child streams without preventing parent records from being emitted when children hit their cap. Bug Fixes: - Apply the dry-run record limit to all streams, including parents and children, instead of only non-child streams. - Prevent child stream aborts due to dry-run record limits from stopping sibling or parent stream processing by catching abort exceptions during child sync. Tests: - Add a regression test verifying that parent and sibling records are still emitted, and child records are correctly capped, when a child stream hits the dry-run record limit. --------- Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
Signed-off-by: Edgar Ramírez-Mondragón <edgarrm358@gmail.com>
8f1375b
into
feat/safely-ignore-errors
Stub
Summary by Sourcery
Track per-stream sync results and allow tap runs to continue syncing remaining streams after individual failures while reporting an aggregate outcome.
New Features:
Bug Fixes:
Enhancements:
Tests: