aminghadersohi commented on code in PR #44009:
URL: https://github.com/apache/superset/pull/44009#discussion_r4007602683


##########
superset/versioning/changes/listener.py:
##########
@@ -395,7 +396,90 @@ def _persist_buffered_records(
         incr_capture_error("bulk_insert")
 
 
-def register_change_record_listener() -> None:  # noqa: C901
+def finalize_change_records(session: Session) -> None:
+    """Build and persist the transaction's change records at commit time.
+
+    Module-level (rather than a closure inside the registration function)
+    so the capture write path can be exercised directly by unit tests
+    against an isolated session; it depends only on the session and the
+    module helpers, never on the registered entity classes.
+    """
+    if session.in_nested_transaction() or session.info.get(_FINALIZING_KEY):
+        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.
+    start: float | None = None
+    try:
+        session.flush()
+        start = perf_counter()
+        initial_states: dict[tuple[str, int], tuple[Any, dict[str, Any]]] = (
+            session.info.get(_INITIAL_STATES_KEY, {})
+        )
+        buffer = _build_scalar_buffer(initial_states)
+
+        try:
+            tx_id = _current_transaction_id(session)
+        except Exception:  # pylint: disable=broad-except
+            logger.exception("version_changes: transaction lookup failed")
+            incr_capture_error("transaction_lookup")
+            return
+        if tx_id is None:
+            return
+
+        _stamp_action_kind_on_transaction(session, tx_id)
+        _append_child_records_to_buffer(session, tx_id, buffer)
+        _inject_action_meta_record(session, buffer)
+
+        if buffer:
+            _persist_buffered_records(session, tx_id, buffer)
+    finally:
+        session.info.pop(_FINALIZING_KEY, None)
+        if start is not None:
+            emit_capture_timing("finalize", (perf_counter() - start) * 1000.0)
+
+
+def _capture_initial_states(
+    session: Session, versioned_classes: tuple[type, ...]
+) -> None:
+    """The ``before_flush`` capture stage: retain each dirty versioned entity's
+    pre-flush database state for the final diff.
+
+    Timed as its own metric stage. The per-entity pre-state SELECTs issued
+    here are the capture cost that scales with the number of dirty versioned
+    entities — on a bulk edit plausibly the dominant cost the kill-switch
+    removes — and they run before the flush, outside ``finalize``'s timer. A
+    sample is emitted only when at least one versioned entity was captured,
+    so the many unrelated autoflushes do not flood the series with empty
+    samples; together with ``finalize`` the two stages cover the whole
+    listener. Module-level (not the registered closure) so it is
+    unit-testable without ``db.session``.
+    """
+    initial_states: dict[tuple[str, int], tuple[Any, dict[str, Any]]] = (
+        session.info.setdefault(_INITIAL_STATES_KEY, {})
+    )
+    start = perf_counter()
+    captured = 0
+    try:
+        for obj in list(session.dirty):
+            if isinstance(obj, versioned_classes):
+                _capture_dirty_entity_initial_state(session, obj, 
initial_states)
+                captured += 1

Review Comment:
   `captured` counts candidates, not captures: 
`_capture_dirty_entity_initial_state` returns early when the key is already in 
`initial_states`. Probe: one entity over 4 flushes issued 1 pre-state SELECT 
but emitted 4 samples, 3 near-zero — diluting the upper percentiles the 
docstring says to alert on.
   
   ```suggestion
           for obj in list(session.dirty):
               if isinstance(obj, versioned_classes):
                   before = len(initial_states)
                   _capture_dirty_entity_initial_state(session, obj, 
initial_states)
                   captured += len(initial_states) - before
   ```
   



##########
superset/versioning/changes/listener.py:
##########
@@ -395,7 +396,90 @@ def _persist_buffered_records(
         incr_capture_error("bulk_insert")
 
 
-def register_change_record_listener() -> None:  # noqa: C901
+def finalize_change_records(session: Session) -> None:
+    """Build and persist the transaction's change records at commit time.
+
+    Module-level (rather than a closure inside the registration function)
+    so the capture write path can be exercised directly by unit tests
+    against an isolated session; it depends only on the session and the
+    module helpers, never on the registered entity classes.
+    """
+    if session.in_nested_transaction() or session.info.get(_FINALIZING_KEY):
+        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.
+    start: float | None = None
+    try:
+        session.flush()
+        start = perf_counter()
+        initial_states: dict[tuple[str, int], tuple[Any, dict[str, Any]]] = (
+            session.info.get(_INITIAL_STATES_KEY, {})
+        )
+        buffer = _build_scalar_buffer(initial_states)
+
+        try:
+            tx_id = _current_transaction_id(session)
+        except Exception:  # pylint: disable=broad-except
+            logger.exception("version_changes: transaction lookup failed")
+            incr_capture_error("transaction_lookup")
+            return
+        if tx_id is None:
+            return
+
+        _stamp_action_kind_on_transaction(session, tx_id)
+        _append_child_records_to_buffer(session, tx_id, buffer)
+        _inject_action_meta_record(session, buffer)
+
+        if buffer:
+            _persist_buffered_records(session, tx_id, buffer)
+    finally:
+        session.info.pop(_FINALIZING_KEY, None)
+        if start is not None:
+            emit_capture_timing("finalize", (perf_counter() - start) * 1000.0)
+
+
+def _capture_initial_states(
+    session: Session, versioned_classes: tuple[type, ...]
+) -> None:
+    """The ``before_flush`` capture stage: retain each dirty versioned entity's
+    pre-flush database state for the final diff.
+
+    Timed as its own metric stage. The per-entity pre-state SELECTs issued
+    here are the capture cost that scales with the number of dirty versioned
+    entities — on a bulk edit plausibly the dominant cost the kill-switch
+    removes — and they run before the flush, outside ``finalize``'s timer. A
+    sample is emitted only when at least one versioned entity was captured,
+    so the many unrelated autoflushes do not flood the series with empty
+    samples; together with ``finalize`` the two stages cover the whole
+    listener. Module-level (not the registered closure) so it is
+    unit-testable without ``db.session``.
+    """
+    initial_states: dict[tuple[str, int], tuple[Any, dict[str, Any]]] = (
+        session.info.setdefault(_INITIAL_STATES_KEY, {})
+    )
+    start = perf_counter()
+    captured = 0
+    try:
+        for obj in list(session.dirty):
+            if isinstance(obj, versioned_classes):
+                _capture_dirty_entity_initial_state(session, obj, 
initial_states)
+                captured += 1
+    finally:

Review Comment:
   `try:`/`finally:` with no `except`: a raise in the per-entity capture 
propagates out of `before_flush` and fails the save (probed). Net-zero vs 
master — the old closure was unguarded too — but it is the twin of the guard 
just added at line 430.
   
   ```suggestion
       except Exception:  # pylint: disable=broad-except
           logger.exception("version_changes: initial-state capture failed")
           incr_capture_error("capture_initial_states")
       finally:
   ```
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to