From 47b79b2b8eb2eca17732e9f66374d2793e3601e5 Mon Sep 17 00:00:00 2001 From: Mike Bridge Date: Fri, 2 Oct 2026 15:35:04 -0600 Subject: [PATCH 1/3] perf(versioning): defer capture policy until versioned work Avoid host policy lookups for unrelated flushes without freezing a synthetic denial. Reuse the connection-scoped decision for baseline capture, semantic changes, and finalization so later flushes cannot disagree. Add native-listener counting and lifecycle regressions for SC-124084. Retain the inherited savepoint rollback cases as strict xfails under the approved Option B scope; master and candidate row-level controls match. Align capture timing documentation with the skipped unrelated commits. --- superset/versioning/baseline/listener.py | 6 +- superset/versioning/changes/listener.py | 16 +- superset/versioning/metrics.py | 6 +- superset/versioning/unit_of_work.py | 43 ++++- tests/unit_tests/versioning/test_listener.py | 16 +- .../versioning/test_runtime_capture.py | 170 ++++++++++++++++++ 6 files changed, 230 insertions(+), 27 deletions(-) diff --git a/superset/versioning/baseline/listener.py b/superset/versioning/baseline/listener.py index 4c22160a205b..c8c57a3de4c5 100644 --- a/superset/versioning/baseline/listener.py +++ b/superset/versioning/baseline/listener.py @@ -49,7 +49,7 @@ ) from superset.versioning.baseline.dirty import force_parent_dirty_on_child_change from superset.versioning.baseline.insertion import insert_baseline_and_children -from superset.versioning.utils import capture_enabled +from superset.versioning.unit_of_work import capture_for_write logger = logging.getLogger(__name__) @@ -102,7 +102,9 @@ def capture_baseline(session: Session, flush_context: Any, instances: Any) -> No # its own ``version_transaction`` row via direct SQL — so without this # guard a detached/kill-switched session would still write baselines. # ``_remove_continuum_write_listeners`` flips this option off. - if not versioning_manager.options["versioning"] or not capture_enabled(session): + if not versioning_manager.options["versioning"] or not capture_for_write( + session + ): return try: # Make sure a child-only edit promotes the parent to diff --git a/superset/versioning/changes/listener.py b/superset/versioning/changes/listener.py index 95229d43ba64..1bf888600e32 100644 --- a/superset/versioning/changes/listener.py +++ b/superset/versioning/changes/listener.py @@ -69,13 +69,13 @@ ) from superset.versioning.metrics import emit_capture_timing, incr_capture_error from superset.versioning.snapshot import reconcile_parent_snapshots -from superset.versioning.utils import capture_enabled +from superset.versioning.unit_of_work import capture_for_write, INITIAL_STATES_KEY logger = logging.getLogger(__name__) # Keys for transaction-scoped state stored on ``session.info``. -_INITIAL_STATES_KEY = "_version_changes_initial_states" +_INITIAL_STATES_KEY: str = INITIAL_STATES_KEY _FINALIZING_KEY = "_version_changes_finalizing" # Key on ``session.info`` that commands set to declare the high-level @@ -447,20 +447,18 @@ def finalize_change_records(session: Session) -> None: against an isolated session; it depends only on the session and the module helpers, never on the registered entity classes. """ - if not capture_enabled(session): - return if session.in_nested_transaction() or session.info.get(_FINALIZING_KEY): return + if not capture_for_write(session): + return session.info[_FINALIZING_KEY] = True # Measures the FINALIZE stage only: the timer starts after the flush, # which excludes the transaction's own write cost but also excludes # capture_initial_states' per-entity pre-state SELECTs (those are timed # as their own ``capture_initial_states`` stage in before_flush) — and - # runs through every capture step and early return. Every commit on the - # session emits a sample, including commits touching no versioned - # entity, because the whole-listener overhead is exactly what the - # kill-switch removes; a flush that raises emits nothing. + # runs through every capture step and early return for allowed versioned + # work. Unrelated commits and a flush that raises emit nothing. start: float | None = None try: session.flush() @@ -582,7 +580,7 @@ def register_change_record_listener() -> None: def capture_initial_states( session: Session, _flush_context: Any, _instances: Any ) -> None: - if not capture_enabled(session): + if not capture_for_write(session): return _capture_initial_states(session, versioned_classes) diff --git a/superset/versioning/metrics.py b/superset/versioning/metrics.py index 1ace45cb92c5..5a46d9deb81a 100644 --- a/superset/versioning/metrics.py +++ b/superset/versioning/metrics.py @@ -67,8 +67,8 @@ def emit_capture_timing(stage: str, duration_ms: float) -> None: edit the dominant cost — sampled whenever at least one pre-state read was attempted, including reads that fail and retain nothing) and ``finalize`` (the post-flush record build and persist, sampled on - every commit on the session, including commits touching no versioned - entity, which still pay the listener overhead). Alert on both, on upper + commits with allowed versioned work; unrelated commits skip capture + and emit no sample). Alert on both, on upper percentiles rather than the mean. :func:`incr_capture_error` covers *loss*; this covers *slowdown*. Best-effort under the same fail-open posture: metrics emission must never itself break a user's save. @@ -81,7 +81,7 @@ def emit_capture_timing(stage: str, duration_ms: float) -> None: f"{_CAPTURE_METRIC_PREFIX}.{stage}.latency", duration_ms ) except Exception as ex: # pylint: disable=broad-except - # This runs on every commit, so a structurally broken stats backend + # This runs on captured commits, so a structurally broken stats backend # (a custom StatsLogger without ``timing()``, or a not-yet-configured # instance at startup) would otherwise log a full traceback per # commit — identical each time. One warning line per occurrence, no diff --git a/superset/versioning/unit_of_work.py b/superset/versioning/unit_of_work.py index 9bcd208fa118..4e128892612d 100644 --- a/superset/versioning/unit_of_work.py +++ b/superset/versioning/unit_of_work.py @@ -16,12 +16,39 @@ # under the License. """Runtime capture suppression at Continuum's transaction-scoped write boundary.""" +from itertools import chain + +from sqlalchemy.engine import Connection from sqlalchemy.orm import Session +from sqlalchemy_continuum import versioning_manager from sqlalchemy_continuum.operation import Operations from sqlalchemy_continuum.unit_of_work import UnitOfWork +from sqlalchemy_continuum.utils import is_versioned from superset.versioning.utils import capture_enabled +INITIAL_STATES_KEY: str = "_version_changes_initial_states" + + +def _has_versioned_work(session: Session) -> bool: + """Recognize pending parents, children, deletes, or retained pre-flush state.""" + return bool(session.info.get(INITIAL_STATES_KEY)) or any( + is_versioned(obj) for obj in chain(session.new, session.dirty, session.deleted) + ) + + +def capture_for_write(session: Session) -> bool: + """Share the unit's frozen decision without evaluating unrelated writes.""" + connection: Connection | None = versioning_manager.session_connection_map.get( + session + ) + unit: CaptureUnitOfWork | None = versioning_manager.units_of_work.get(connection) + if unit is None: + if not _has_versioned_work(session): + return False + unit = versioning_manager.unit_of_work(session) + return unit.capture_if_needed(session) + class CaptureUnitOfWork(UnitOfWork): """Keep denied mapper/association operations out of later captured flushes.""" @@ -37,9 +64,21 @@ def _capture_enabled(self, session: Session) -> bool: self._capture_allowed = capture_enabled(session) return self._capture_allowed + def capture_if_needed(self, session: Session) -> bool: + """Leave unrelated work undecided until the first versioned flush.""" + if session is self.version_session: + return False + if ( + self._capture_allowed is None + and not self.has_changes + and not _has_versioned_work(session) + ): + return False + return self._capture_enabled(session) + def process_before_flush(self, session: Session) -> None: """Decide before Continuum creates its transaction or version session.""" - if session is self.version_session or not self._capture_enabled(session): + if not self.capture_if_needed(session): return super().process_before_flush(session) @@ -47,7 +86,7 @@ def process_after_flush(self, session: Session) -> None: """Discard denied operations, including relationship-table statements.""" if session is self.version_session: return - if not self._capture_enabled(session): + if not self.capture_if_needed(session): self.operations: Operations = Operations() self.pending_statements.clear() return diff --git a/tests/unit_tests/versioning/test_listener.py b/tests/unit_tests/versioning/test_listener.py index d14fbab0a69d..cef610b1e0ce 100644 --- a/tests/unit_tests/versioning/test_listener.py +++ b/tests/unit_tests/versioning/test_listener.py @@ -430,15 +430,10 @@ def test_transient_persist_failure_is_logged_and_counted( metric_spy.assert_called_once_with("bulk_insert") -def test_capture_latency_metric_fires_on_commit( +def test_capture_latency_metric_skips_nonversioned_commit( lifecycle_session: Session, mocker: Any ) -> None: - """The finalizer emits the write-path latency series on every save-path - commit — the kill-switch's own decision signal, measuring capture - overhead only (the timer starts after the transaction's own flush). - Driven through the real module-level finalizer on an isolated session - (no versioning tables needed: the tx-id early return still passes the - timing's ``finally``).""" + """Unrelated commits do not dilute the version-capture latency series.""" sa.event.listen( lifecycle_session, "before_commit", listener.finalize_change_records ) @@ -452,10 +447,7 @@ def test_capture_latency_metric_fires_on_commit( for call in manager.instance.timing.call_args_list if call.args[0] == "superset.versioning.capture.finalize.latency" ] - assert len(calls) == 1 - duration_ms: float = calls[0].args[1] - assert isinstance(duration_ms, float) - assert duration_ms >= 0 + assert calls == [] def test_capture_latency_metric_skips_reentrant_finalize( @@ -493,6 +485,7 @@ def test_capture_latency_metric_emits_nothing_when_flush_fails( """A flush that raises is the user's own failing write, not capture cost: the exception propagates and no sample lands in the series.""" manager: MagicMock = MagicMock() + mocker.patch.object(listener, "capture_for_write", return_value=True) mocker.patch("superset.extensions.stats_logger_manager", manager) session: MagicMock = MagicMock() session.info = {} @@ -602,6 +595,7 @@ def test_transaction_lookup_failure_does_not_break_the_commit( lifecycle_session, "before_commit", listener.finalize_change_records ) mocker.patch("superset.extensions.stats_logger_manager", MagicMock()) + mocker.patch.object(listener, "capture_for_write", return_value=True) error_spy: MagicMock = mocker.patch.object(listener, "incr_capture_error") def explode(session: Session) -> int: diff --git a/tests/unit_tests/versioning/test_runtime_capture.py b/tests/unit_tests/versioning/test_runtime_capture.py index 043c6fbd82dd..46dffce39176 100644 --- a/tests/unit_tests/versioning/test_runtime_capture.py +++ b/tests/unit_tests/versioning/test_runtime_capture.py @@ -17,6 +17,7 @@ """Persisted saves through the real baseline, Continuum and change listeners.""" from collections.abc import Iterator +from datetime import datetime, timezone from itertools import chain, repeat from typing import Any from unittest.mock import MagicMock, patch @@ -394,3 +395,172 @@ def test_none_predicate_result_denies_capture_and_is_memoized( assert unit._capture_enabled(capture_session) is False assert unit._capture_enabled(capture_session) is False predicate.assert_called_once_with(capture_session) + + +@pytest.mark.parametrize("enabled", [False, True]) +def test_nonversioned_transaction_never_consults_capture_policy( + capture_session: Session, + app: SupersetApp, + monkeypatch: pytest.MonkeyPatch, + enabled: bool, +) -> None: + """Unrelated writes and empty commits do not ask the host for policy.""" + predicate: MagicMock = MagicMock(return_value=enabled) + monkeypatch.setitem(app.config, "VERSIONING_CAPTURE_PREDICATE", predicate) + database: Database = Database(database_name="unrelated", sqlalchemy_uri="sqlite://") + capture_session.add(database) + capture_session.flush() + database.database_name = "edited" + capture_session.commit() + capture_session.commit() + assert not any(history_counts(capture_session).values()) + predicate.assert_not_called() + + +@pytest.mark.parametrize("unrelated_first", [False, True]) +@pytest.mark.parametrize("enabled", [False, True]) +def test_first_versioned_flush_freezes_policy_through_finalization( + capture_session: Session, + app: SupersetApp, + monkeypatch: pytest.MonkeyPatch, + unrelated_first: bool, + enabled: bool, +) -> None: + """One decision covers shadows and semantic changes across mixed flushes.""" + dashboard: Dashboard = Dashboard(dashboard_title="original") + capture_session.add(dashboard) + capture_session.commit() + before: dict[str, int] = history_counts(capture_session) + predicate: MagicMock = MagicMock(return_value=enabled) + monkeypatch.setitem(app.config, "VERSIONING_CAPTURE_PREDICATE", predicate) + if unrelated_first: + capture_session.add( + Database(database_name="unrelated", sqlalchemy_uri="sqlite://") + ) + capture_session.flush() + predicate.assert_not_called() + dashboard.dashboard_title = "intermediate" + capture_session.flush() + predicate.assert_called_once_with(capture_session) + predicate.return_value = not enabled + dashboard.dashboard_title = "final" + capture_session.commit() + predicate.assert_called_once_with(capture_session) + after: dict[str, int] = history_counts(capture_session) + assert after["dashboards_version"] == before["dashboards_version"] + int(enabled) + assert (after["version_changes"] > before["version_changes"]) is enabled + + +@pytest.mark.parametrize( + "enabled,rollback_nested", + [ + (False, False), + (False, True), + (True, False), + pytest.param( + True, + True, + marks=pytest.mark.xfail( + strict=True, + raises=sa.exc.OperationalError, + reason="SC-TBD: inherited Continuum savepoint rollback lifecycle bug", + ), + ), + ], +) +def test_lazy_capture_savepoint_and_query_autoflush( + capture_session: Session, + app: SupersetApp, + monkeypatch: pytest.MonkeyPatch, + rollback_nested: bool, + enabled: bool, +) -> None: + """Savepoint completion keeps the outer decision through query autoflush.""" + predicate: MagicMock = MagicMock(return_value=enabled) + monkeypatch.setitem(app.config, "VERSIONING_CAPTURE_PREDICATE", predicate) + capture_session.add(Database(database_name="unrelated", sqlalchemy_uri="sqlite://")) + nested: SessionTransaction = capture_session.begin_nested() + predicate.assert_not_called() + dashboard: Dashboard = Dashboard(dashboard_title="nested") + capture_session.add(dashboard) + assert ( + capture_session.scalar(sa.select(sa.func.count()).select_from(Dashboard)) == 1 + ) + predicate.assert_called_once_with(capture_session) + predicate.return_value = not enabled + if rollback_nested: + nested.rollback() + else: + nested.commit() + capture_session.add(Dashboard(dashboard_title="outer")) + capture_session.commit() + predicate.assert_called_once_with(capture_session) + assert history_counts(capture_session)["dashboards_version"] == ( + (1 if rollback_nested else 2) if enabled else 0 + ) + + +@pytest.mark.parametrize("soft_delete", [False, True]) +@pytest.mark.parametrize("enabled", [False, True]) +def test_lazy_capture_delete_policy( + capture_session: Session, + app: SupersetApp, + monkeypatch: pytest.MonkeyPatch, + soft_delete: bool, + enabled: bool, +) -> None: + """Both deletion paths consult policy once and preserve allowed history.""" + dashboard: Dashboard = Dashboard(dashboard_title="delete me") + capture_session.add(dashboard) + capture_session.commit() + before: dict[str, int] = history_counts(capture_session) + predicate: MagicMock = MagicMock(return_value=enabled) + monkeypatch.setitem(app.config, "VERSIONING_CAPTURE_PREDICATE", predicate) + if soft_delete: + dashboard.deleted_at = datetime.now(timezone.utc) + else: + capture_session.delete(dashboard) + capture_session.commit() + predicate.assert_called_once_with(capture_session) + assert history_counts(capture_session)["dashboards_version"] == ( + before["dashboards_version"] + int(enabled and not soft_delete) + ) + + +def test_lazy_capture_predicate_error_still_propagates( + capture_session: Session, + app: SupersetApp, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Laziness does not turn a host programming failure into silent denial.""" + predicate: MagicMock = MagicMock(side_effect=RuntimeError("policy failed")) + monkeypatch.setitem(app.config, "VERSIONING_CAPTURE_PREDICATE", predicate) + capture_session.add(Dashboard(dashboard_title="not committed")) + with pytest.raises(RuntimeError, match="policy failed"): + capture_session.commit() + capture_session.rollback() + assert not any(history_counts(capture_session).values()) + + +@pytest.mark.xfail( + strict=True, + raises=sa.exc.OperationalError, + reason="SC-TBD: inherited Continuum savepoint rollback lifecycle bug", +) +def test_capture_after_savepoint_rollback_with_stable_policy( + capture_session: Session, + app: SupersetApp, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A rolled-back shadow cannot poison a subsequent outer-transaction save.""" + monkeypatch.setitem( + app.config, "VERSIONING_CAPTURE_PREDICATE", lambda session: True + ) + capture_session.add(Database(database_name="unrelated", sqlalchemy_uri="sqlite://")) + nested: SessionTransaction = capture_session.begin_nested() + capture_session.add(Dashboard(dashboard_title="rolled back")) + capture_session.flush() + nested.rollback() + capture_session.add(Dashboard(dashboard_title="outer")) + capture_session.commit() + assert history_counts(capture_session)["dashboards_version"] == 1 From a199a7728ddb87bad38d4845b3e689e3547171c0 Mon Sep 17 00:00:00 2001 From: Mike Bridge Date: Fri, 2 Oct 2026 16:20:24 -0600 Subject: [PATCH 2/3] test(versioning): link savepoint xfails to SC-124128 Replace the provisional ticket references with the filed savepoint lifecycle bug so the inherited strict xfails remain traceable. Capture logic and xfail constraints are unchanged. --- tests/unit_tests/versioning/test_runtime_capture.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/unit_tests/versioning/test_runtime_capture.py b/tests/unit_tests/versioning/test_runtime_capture.py index 46dffce39176..394b6e7f38f8 100644 --- a/tests/unit_tests/versioning/test_runtime_capture.py +++ b/tests/unit_tests/versioning/test_runtime_capture.py @@ -463,7 +463,7 @@ def test_first_versioned_flush_freezes_policy_through_finalization( marks=pytest.mark.xfail( strict=True, raises=sa.exc.OperationalError, - reason="SC-TBD: inherited Continuum savepoint rollback lifecycle bug", + reason="SC-124128: inherited Continuum savepoint rollback bug", ), ), ], @@ -545,7 +545,7 @@ def test_lazy_capture_predicate_error_still_propagates( @pytest.mark.xfail( strict=True, raises=sa.exc.OperationalError, - reason="SC-TBD: inherited Continuum savepoint rollback lifecycle bug", + reason="SC-124128: inherited Continuum savepoint rollback bug", ) def test_capture_after_savepoint_rollback_with_stable_policy( capture_session: Session, From d2d69c36327717d9a98492689e3e5a54f537d5a0 Mon Sep 17 00:00:00 2001 From: Mike Bridge Date: Thu, 8 Oct 2026 09:43:02 -0600 Subject: [PATCH 3/3] test(versioning): remove private tracker references from xfail reasons --- tests/unit_tests/versioning/test_runtime_capture.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/unit_tests/versioning/test_runtime_capture.py b/tests/unit_tests/versioning/test_runtime_capture.py index 394b6e7f38f8..844fc27dd4c2 100644 --- a/tests/unit_tests/versioning/test_runtime_capture.py +++ b/tests/unit_tests/versioning/test_runtime_capture.py @@ -463,7 +463,7 @@ def test_first_versioned_flush_freezes_policy_through_finalization( marks=pytest.mark.xfail( strict=True, raises=sa.exc.OperationalError, - reason="SC-124128: inherited Continuum savepoint rollback bug", + reason="Inherited Continuum savepoint rollback bug", ), ), ], @@ -545,7 +545,7 @@ def test_lazy_capture_predicate_error_still_propagates( @pytest.mark.xfail( strict=True, raises=sa.exc.OperationalError, - reason="SC-124128: inherited Continuum savepoint rollback bug", + reason="Inherited Continuum savepoint rollback bug", ) def test_capture_after_savepoint_rollback_with_stable_policy( capture_session: Session,