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]
