mikebridge commented on code in PR #45029:
URL: https://github.com/apache/superset/pull/45029#discussion_r4215826845


##########
UPDATING.md:
##########
@@ -1709,6 +1709,8 @@ Authorization reuses the resource's `can_read` permission 
and per-object `raise_
 
 Entity version history (the `version_transaction` / `*_version` shadow tables 
that back version capture) is aged out by a nightly Celery beat task, 
`version_history.prune_old_versions` 
(`superset.tasks.version_history_retention`).
 
+**Scheduled-run cap (behavior change):** 
`VERSION_HISTORY_PRUNE_MAX_TRANSACTIONS_PER_RUN` defaults to 1000 whole 
prunable transactions per invocation; associated shadow and change rows remain 
atomic. `0` or `None` restores unlimited runs. Invalid values skip the task 
before deletion. Capped results report `cap_reached`, `remaining_eligible`, and 
`remaining_count_complete`; the remainder is a lower bound from one candidate 
window after the live scan's stopping point when the probe does not cover the 
backlog, and `None` if the probe fails after committed deletion. Physical 
shadow-row counts remain separate. `VERSION_HISTORY_PRUNE_DRY_RUN=True` scans 
the entire prunable backlog without writes and reports the exact 
`eligible_backlog` and `estimated_capped_runs`; non-boolean values skip the 
scheduled task without deleting history. A large backlog now drains over 
successive scheduled runs rather than in one invocation. The cap is per task 
invocation, not a quota shared across concurrent w
 orkers. If the database commits a window but the client loses its 
acknowledgement, retries can prune additional windows beyond the cap (up to two 
extra windows per retried pass with the three-attempt policy); all pruned 
transactions still meet the retention policy.

Review Comment:
   Thanks, the old sentence was stale. In 
[7868f04f16](https://github.com/apache/superset/commit/7868f04f16fa42625f95ecbc00f410412ede2569),
 `UPDATING.md` now says that when a commit outcome is classified as uncertain, 
the task reports an error and stops without retrying that window, so that 
outcome cannot prune past the cap. I scoped the wording to the classification 
the code can establish; the retry path and task error handling support it.



##########
tests/unit_tests/tasks/test_version_history_retention.py:
##########
@@ -407,9 +410,532 @@ def test_immediate_cutoff_and_invalid_skip(stats: 
MagicMock, days: int) -> None:
             days
         )
     if days == -1:
-        run_pass.assert_called_once_with(now, [], 0)
+        run_pass.assert_called_once_with(now, [], 0, max_prune=None)
         assert result["cutoff"] == now.isoformat()
     else:
         assert result == {"skipped": 1}
         tables.assert_not_called()
         run_pass.assert_not_called()
+
+
+def test_prune_cap_accounts_only_for_committed_passes(stats: MagicMock) -> 
None:
+    """A retried pass uses the same allowance and cannot overshoot the cap."""
+    tables: version_history_retention.ShadowTables = (
+        version_history_retention.ShadowTables(
+            parent=[], child=[], m2m=None, transaction=MagicMock()
+        )
+    )
+    first: dict[str, int] = {
+        "candidate_count": 1000,
+        "max_candidate_id": 1000,
+        "pruned_transactions": 2,
+    }
+    second: dict[str, int] = {
+        "candidate_count": 1000,
+        "max_candidate_id": 2000,
+        "pruned_transactions": 1,
+    }
+    run_pass: MagicMock
+    with (
+        patch.object(
+            version_history_retention, "_resolve_shadow_tables", 
return_value=tables
+        ),
+        patch.object(
+            version_history_retention,
+            "_run_pass_with_retry",
+            side_effect=[(first, 1), (second, 0)],
+        ) as run_pass,
+        patch.object(
+            version_history_retention, "_probe_prunable", return_value=(4, 
False)
+        ),
+    ):
+        result: dict[str, Any] = 
version_history_retention._prune_old_versions_impl(
+            retention_days=30, max_per_run=3
+        )
+
+    assert result["pruned_transactions"] == 3
+    assert result["retried"] == 1
+    assert result["cap_reached"] is True
+    assert result["remaining_eligible"] == 4
+    assert result["remaining_count_complete"] is False
+    stats.gauge.assert_any_call(
+        "superset.versioning.retention.remaining_count_complete", 0
+    )
+    assert run_pass.call_args_list[0].kwargs["max_prune"] == 3
+    assert run_pass.call_args_list[1].kwargs["max_prune"] == 1
+
+
+def test_prune_cap_probes_after_preserved_candidate_windows(
+    stats: MagicMock,
+) -> None:
+    """The remainder probe skips windows already scanned before the cap."""
+    tables: version_history_retention.ShadowTables = (
+        version_history_retention.ShadowTables(
+            parent=[], child=[], m2m=None, transaction=MagicMock()
+        )
+    )
+    batch_size: int = version_history_retention._MAX_PRUNE_BATCH
+    windows: dict[int, version_history_retention._PruneWindow] = {
+        0: version_history_retention._PruneWindow([], batch_size, 1000),
+        1000: version_history_retention._PruneWindow([1500], batch_size, 2000),
+    }
+
+    def resolve_window(
+        _conn: sa.engine.Connection,
+        _cutoff: datetime,
+        _shadow_tables: list[sa.Table],
+        after_id: int,
+        _limit: int,
+    ) -> version_history_retention._PruneWindow:
+        """Return the candidate window at the requested scan cursor."""
+        return windows[after_id]
+
+    resolve: MagicMock
+    with (
+        patch.object(
+            version_history_retention, "_resolve_shadow_tables", 
return_value=tables
+        ),
+        patch.object(
+            version_history_retention,
+            "_run_pass_with_retry",
+            side_effect=[
+                ({"candidate_count": batch_size, "max_candidate_id": 1000}, 0),
+                (
+                    {
+                        "candidate_count": batch_size,
+                        "max_candidate_id": 2000,
+                        "pruned_transactions": 1,
+                    },
+                    0,
+                ),
+            ],
+        ),
+        patch.object(version_history_retention, "db"),
+        patch.object(
+            version_history_retention,
+            "_resolve_prune_window",
+            side_effect=resolve_window,
+        ) as resolve,
+    ):
+        result: dict[str, Any] = 
version_history_retention._prune_old_versions_impl(
+            retention_days=30, max_per_run=1
+        )
+
+    assert result["pruned_transactions"] == 1
+    assert result["remaining_eligible"] == 1
+    assert result["remaining_count_complete"] is False
+    assert resolve.call_args.args[3] == 1000
+    stats.gauge.assert_any_call(
+        "superset.versioning.retention.remaining_eligible_at_least", 1
+    )
+
+
[email protected]("cap", [None, 0, 3])
+def test_prune_dry_run_counts_all_eligible_without_writes(
+    stats: MagicMock, cap: int | None
+) -> None:
+    """A dry run reports the whole backlog and capped-run estimate."""
+    tables: version_history_retention.ShadowTables = (
+        version_history_retention.ShadowTables(
+            parent=[], child=[], m2m=None, transaction=MagicMock()
+        )
+    )
+    run_pass: MagicMock
+    with (
+        patch.object(
+            version_history_retention, "_resolve_shadow_tables", 
return_value=tables
+        ),
+        patch.object(version_history_retention, "_count_prunable", 
return_value=7),
+        patch.object(version_history_retention, "_run_pass_with_retry") as 
run_pass,
+    ):
+        result: dict[str, Any] = 
version_history_retention._prune_old_versions_impl(
+            30, max_per_run=cap, dry_run=True
+        )
+
+    assert result["eligible_backlog"] == 7
+    assert result["estimated_capped_runs"] == (3 if cap == 3 else 1)
+    run_pass.assert_not_called()
+    stats.gauge.assert_not_called()
+
+
[email protected]("helper_name", ["_count_prunable", "_probe_prunable"])
+def test_remainder_reads_use_default_isolation(helper_name: str) -> None:
+    """Backlog measurement does not request the delete pass's isolation."""
+    tables: version_history_retention.ShadowTables = (
+        version_history_retention.ShadowTables(
+            parent=[], child=[], m2m=None, transaction=MagicMock()
+        )
+    )
+    engine: MagicMock = MagicMock()
+    window: version_history_retention._PruneWindow = (
+        version_history_retention._PruneWindow(
+            prunable=[], candidate_count=0, max_candidate_id=0
+        )
+    )
+    with (
+        patch.object(version_history_retention, "db") as mock_db,
+        patch.object(
+            version_history_retention, "_resolve_prune_window", 
return_value=window
+        ),
+    ):
+        mock_db.engine = engine
+        result: int | tuple[int, bool] = getattr(
+            version_history_retention, helper_name
+        )(datetime(2026, 1, 1), tables)
+
+    assert result in (0, (0, True))
+    engine.connect.return_value.execution_options.assert_not_called()
+
+
+def test_dry_run_count_accumulates_windows_and_closes_read_transactions() -> 
None:
+    """A paged backlog count sums windows and releases each read snapshot."""
+    tables: version_history_retention.ShadowTables = (
+        version_history_retention.ShadowTables(
+            parent=[], child=[], m2m=None, transaction=MagicMock()
+        )
+    )
+    engine: MagicMock = MagicMock()
+    full_window: version_history_retention._PruneWindow = (
+        version_history_retention._PruneWindow(
+            prunable=[1, 2],
+            candidate_count=version_history_retention._MAX_PRUNE_BATCH,
+            max_candidate_id=1000,
+        )
+    )
+    final_window: version_history_retention._PruneWindow = (
+        version_history_retention._PruneWindow(
+            prunable=[1001], candidate_count=1, max_candidate_id=1001
+        )
+    )
+    mock_db: MagicMock
+    resolve: MagicMock
+    with (
+        patch.object(version_history_retention, "db") as mock_db,
+        patch.object(
+            version_history_retention,
+            "_resolve_prune_window",
+            side_effect=[full_window, final_window],
+        ) as resolve,
+    ):
+        mock_db.engine = engine
+        count: int = version_history_retention._count_prunable(
+            datetime(2026, 1, 1), tables
+        )
+
+    assert count == 3
+    assert resolve.call_args_list[0].args[3] == 0
+    assert resolve.call_args_list[1].args[3] == full_window.max_candidate_id
+    assert engine.connect.call_count == 2
+    assert engine.connect.return_value.__exit__.call_count == 2
+    first_exit: int = engine.mock_calls.index(call.connect().__exit__(None, 
None, None))
+    second_connect: int = engine.mock_calls.index(call.connect(), first_exit + 
1)
+    assert first_exit < second_connect
+
+
[email protected]("invalid", [-1, True, "3", 2.5])
+def test_prune_rejects_invalid_cap_before_work(invalid: object) -> None:
+    """Malformed budgets fail closed before resolving any shadow table."""
+    with (
+        patch.object(version_history_retention, "_resolve_shadow_tables") as 
resolve,
+        pytest.raises(ValueError, match="VERSION_HISTORY_PRUNE_MAX"),
+    ):
+        version_history_retention._prune_old_versions_impl(
+            30, max_per_run=cast(int | None, invalid)
+        )
+    resolve.assert_not_called()
+
+
+def test_prune_retry_reuses_the_same_transaction_budget(stats: MagicMock) -> 
None:
+    """A rolled-back SERIALIZABLE attempt consumes no transaction allowance."""
+    tables: version_history_retention.ShadowTables = (
+        version_history_retention.ShadowTables(
+            parent=[], child=[], m2m=None, transaction=MagicMock()
+        )
+    )
+    cutoff: datetime = datetime(2026, 1, 1)
+    run_pass: MagicMock
+    with (
+        patch.object(
+            version_history_retention,
+            "_run_prune_pass",
+            side_effect=[
+                OperationalError("SELECT 1", {}, Exception("serialization 
conflict")),
+                {"pruned_transactions": 2},
+            ],
+        ) as run_pass,
+        patch.object(version_history_retention.time, "sleep"),
+    ):
+        result: tuple[dict[str, Any], int] = (
+            version_history_retention._run_pass_with_retry(
+                cutoff, tables, after_id=4, max_prune=2
+            )
+        )
+
+    assert result == ({"pruned_transactions": 2}, 1)
+    assert run_pass.call_args_list == [
+        call(cutoff, tables, 4, max_prune=2),
+        call(cutoff, tables, 4, max_prune=2),
+    ]
+    stats.incr.assert_called_once_with("superset.versioning.retention.retried")
+
+
[email protected](
+    "commit_error",
+    [
+        OperationalError(
+            "COMMIT", {}, Exception("connection lost"), 
connection_invalidated=True
+        ),
+        RuntimeError("acknowledgement lost"),
+    ],
+)
+def test_prune_does_not_retry_an_uncertain_commit(
+    stats: MagicMock, commit_error: Exception
+) -> None:
+    """A lost commit acknowledgement must not spend another prune window."""
+    tables: version_history_retention.ShadowTables = (
+        version_history_retention.ShadowTables(
+            parent=[],
+            child=[],
+            m2m=None,
+            transaction=sa.table("version_transaction", sa.column("id")),
+        )
+    )
+    window: version_history_retention._PruneWindow = (
+        version_history_retention._PruneWindow(
+            prunable=[1], candidate_count=1, max_candidate_id=1
+        )
+    )
+    engine: MagicMock = MagicMock()
+    engine_connection: MagicMock = engine.connect.return_value
+    connection: MagicMock = (
+        engine_connection.execution_options.return_value.__enter__.return_value
+    )
+    transaction: MagicMock = connection.begin.return_value
+    transaction.commit.side_effect = commit_error
+    mock_db: MagicMock
+    sleep: MagicMock
+    with (
+        patch.object(version_history_retention, "db") as mock_db,
+        patch.object(
+            version_history_retention, "_resolve_prune_window", 
return_value=window
+        ),
+        patch.object(
+            version_history_retention, "_delete_for_transactions", 
return_value=0

Review Comment:
   Agreed. 
[7868f04f16](https://github.com/apache/superset/commit/7868f04f16fa42625f95ecbc00f410412ede2569)
 adds a `_run_prune_pass` case with prunable IDs `[1, 2, 3]` and `max_prune=2`. 
It asserts parent, child, M2M and transaction deletes all receive `[1, 2]`. The 
test fails if the parent shadow delete is changed to use the uncapped window, 
and the full focused file passes (67 tests).



##########
superset/tasks/deletion_retention.py:
##########
@@ -342,6 +492,48 @@ def _finalize_blocked(record_id: UUID | None, blocker: 
BlockerReason) -> None:
         )
 
 
+def _confirm_committed_purge(record_id: UUID | None, result: CascadeResult) -> 
None:
+    """Keep a committed root counted even if audit finalization fails."""
+    try:
+        audit.confirm(
+            record_id,
+            affected_referrers=result.dangling_chart_uuids,
+            removed_dashboard_slices=result.removed_dashboard_slices,
+        )
+    except Exception:  # pylint: disable=broad-except
+        # The audit remains pending for its normal reconciliation path.
+        logger.exception(
+            "deletion_retention: audit finalization failed after committed 
purge"
+        )
+
+
+def _commit_purge_root(record_id: UUID | None) -> None:
+    """Commit a purge, keeping uncertain outcomes pending for 
reconciliation."""
+    try:
+        db.session.commit()  # pylint: disable=consider-using-transaction
+    except Exception as exc:
+        try:
+            db.session.rollback()  # pylint: disable=consider-using-transaction
+        except Exception:  # pylint: disable=broad-except
+            logger.exception("deletion_retention: rollback after commit error 
failed")
+            raise _PurgeCommitUncertainError(

Review Comment:
   Agreed. 
[7868f04f16](https://github.com/apache/superset/commit/7868f04f16fa42625f95ecbc00f410412ede2569)
 adds a failed-commit plus failed-rollback case through `_purge_model` with two 
eligible roots. It asserts the first root's outcome stays uncertain, the audit 
remains pending, and the second root is not reached. Replacing the 
uncertain-error raise with `raise exc` makes this test fail; the focused file 
passes (110 tests).



-- 
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