diff --git a/UPDATING.md b/UPDATING.md index 28617c43c5ca..52d2ff71de26 100644 --- a/UPDATING.md +++ b/UPDATING.md @@ -491,7 +491,7 @@ Purging is **live by default** (`SOFT_DELETE_PURGE_DRY_RUN=False`), so the reten Deployments that replace the default `CELERY_CONFIG` must ensure workers register `superset.tasks.deletion_retention` and schedule the `deletion_retention.purge_soft_deleted` task themselves. The shipped Docker development config uses `imports` and includes both entries. While `SOFT_DELETE` is statically enabled, a missing beat entry logs a startup warning; when the override explicitly defines `imports`, a missing purge module is also reported. -Operators can immediately erase a specific entity for compliance (GDPR) via `superset deletion-retention force-purge --uuid `; this applies legacy hard-delete semantics — a live chart referencing a force-purged dataset is left without a datasource until re-pointed (the chart is not modified), and it purges the named entity even when it was never soft-deleted. Every purge writes an immutable, content-free audit record to the new `purge_audit_log` table that survives the entity it names: the **scheduled** purge fails closed (an entity whose audit row cannot be written is skipped and retried next run), while **force-purge** proceeds even if the audit write fails — the operator is present and deletion outranks audit for a compliance erasure. +Operators can immediately erase a specific entity for compliance (GDPR) via `superset deletion-retention force-purge --uuid `; this applies legacy hard-delete semantics — a live chart referencing a force-purged dataset is left without a datasource until re-pointed (the chart is not modified), and it purges the named entity even when it was never soft-deleted. Every scheduled evaluation writes a provisional, content-free record to the new `purge_audit_log` table before the cascade starts. Meaningful retained outcomes survive the entity they name. Consecutive scheduled evaluations with the same blocked outcome suppress only the redundant current provisional record; completed outcomes, outcome transitions, and every force-purge attempt remain independent and immutable. The **scheduled** purge fails closed when its provisional record cannot be written, while **force-purge** proceeds even if the audit write fails — the operator is present and deletion outranks audit for a compliance erasure. Operators can monitor `deletion_retention.blocked_audit_suppressed` and `deletion_retention.blocked_audit_dedupe_fallback` to verify suppression and fail-safe fallback behavior without changing the existing blocked-workload gauge. ### Recently Archived view and permanent delete (purge) endpoints diff --git a/superset/commands/deletion_retention/audit.py b/superset/commands/deletion_retention/audit.py index 90d2ea8e39c7..ffb3c624f72e 100644 --- a/superset/commands/deletion_retention/audit.py +++ b/superset/commands/deletion_retention/audit.py @@ -14,16 +14,18 @@ # KIND, either express or implied. See the License for the # specific language governing permissions and limitations # under the License. -"""Write-ahead purge audit record. +"""Write-ahead purge audit records and retained purge outcomes. -Every purge — time-based or force — writes an immutable record that -**survives** the entity it names, on a **dedicated session** outside the +Every purge evaluation — scheduled or force — writes a provisional record on +a **dedicated session** outside the purge transaction so it neither entangles with the ``DBEventLogger`` (which shares ``db.session`` and commits mid-request) nor vanishes if the purge rolls back. The record is written ``pending`` *before* the purge and flipped to ``confirmed`` *after* it commits, so a crash leaves at most a ``pending`` row, never a missing one. ``pending`` rows are reconciled on the -next run (the purge is convergent). +next run. Completed records are immutable. Consecutive scheduled evaluations +that remain blocked may discard only their current redundant provisional row; +force-purge and other meaningful outcomes are retained independently. The dedicated ``purge_audit_log`` table is content-free (no name or PII; only action, actor, UTC time, entity type, UUID, and affected referrers) and is never @@ -39,11 +41,13 @@ from __future__ import annotations import logging +from dataclasses import dataclass from datetime import datetime, timedelta, timezone -from typing import Any, cast +from typing import Any, cast, Literal, TypeAlias from uuid import UUID import sqlalchemy as sa +from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.orm import Session, sessionmaker from superset import db @@ -73,6 +77,17 @@ def _dedicated_session() -> Session: ACTOR_SYSTEM = "system" +RetentionBlockedDisposition: TypeAlias = Literal["retained", "suppressed", "fallback"] + + +@dataclass(frozen=True) +class _AuditRecoverySnapshot: + id: UUID + actor: str + entity_type: str + entity_uuid: str | None + created_on: datetime + def _utc_now() -> datetime: """Naive UTC now for the audit columns. @@ -182,6 +197,146 @@ def block(record_id: UUID | None) -> None: finalize(record_id, STATUS_BLOCKED) +def _capture_recovery_snapshot(record: PurgeAuditLog) -> _AuditRecoverySnapshot: + """Capture the content-free fields needed for fail-safe recovery.""" + return _AuditRecoverySnapshot( + id=cast(UUID, record.id), + actor=str(record.actor), + entity_type=str(record.entity_type), + entity_uuid=record.entity_uuid, + created_on=cast(datetime, record.created_on), + ) + + +def _retention_predecessor( + session: Session, current: PurgeAuditLog +) -> PurgeAuditLog | None: + """Return the latest row that could unambiguously precede ``current``.""" + return session.execute( + sa.select(PurgeAuditLog) + .where(PurgeAuditLog.entity_uuid == current.entity_uuid) + .where(PurgeAuditLog.entity_type == current.entity_type) + .where(PurgeAuditLog.trigger == TRIGGER_RETENTION) + .where(PurgeAuditLog.created_on <= current.created_on) + .where(PurgeAuditLog.id != current.id) + .order_by(PurgeAuditLog.created_on.desc()) + .limit(1) + ).scalar_one_or_none() + + +def _suppress_redundant_block( + session: Session, current: PurgeAuditLog, predecessor: PurgeAuditLog | None +) -> bool: + """Delete only a pending row with a strictly older blocked predecessor.""" + if ( + predecessor is None + or predecessor.created_on >= current.created_on + or predecessor.status != STATUS_BLOCKED + ): + return False + deleted_rows: int | None = session.execute( + sa.delete(PurgeAuditLog.__table__).where( + PurgeAuditLog.__table__.c.id == current.id, + PurgeAuditLog.__table__.c.status == STATUS_PENDING, + ) + ).rowcount + if deleted_rows not in (0, 1): + raise SQLAlchemyError( + f"indeterminate purge audit suppression rowcount: {deleted_rows}" + ) + return deleted_rows == 1 + + +def _retain_blocked(session: Session, record_id: UUID) -> None: + """Conditionally retain the current provisional row as blocked.""" + session.execute( + sa.update(PurgeAuditLog.__table__) + .where( + PurgeAuditLog.__table__.c.id == record_id, + PurgeAuditLog.__table__.c.status == STATUS_PENDING, + ) + .values(status=STATUS_BLOCKED, removed_dashboard_slices=0) + ) + + +def _recover_retention_blocked( + record_id: UUID, snapshot: _AuditRecoverySnapshot | None +) -> RetentionBlockedDisposition: + """Retain blocked evidence on a fresh session after persistence uncertainty.""" + recovery_session: Session = _dedicated_session() + try: + current: PurgeAuditLog | None = recovery_session.get(PurgeAuditLog, record_id) + if current is not None: + if current.status == STATUS_PENDING: + _retain_blocked(recovery_session, record_id) + recovery_session.commit() + return "fallback" + if snapshot is None: + return "fallback" + recovery_session.add( + PurgeAuditLog( + id=snapshot.id, + status=STATUS_BLOCKED, + trigger=TRIGGER_RETENTION, + actor=snapshot.actor, + entity_type=snapshot.entity_type, + entity_uuid=snapshot.entity_uuid, + removed_dashboard_slices=0, + created_on=snapshot.created_on, + ) + ) + recovery_session.commit() + except SQLAlchemyError: + recovery_session.rollback() + logger.warning( + "deletion_retention: failed to recover blocked audit row %s " + "entity_type=%s entity_uuid=%s", + record_id, + snapshot.entity_type if snapshot else None, + snapshot.entity_uuid if snapshot else None, + exc_info=True, + ) + finally: + recovery_session.close() + return "fallback" + + +def finalize_retention_blocked( + record_id: UUID | None, +) -> RetentionBlockedDisposition: + """Finalize a scheduled blocker, suppressing only proven redundant evidence.""" + if record_id is None: + return "fallback" + session: Session = _dedicated_session() + snapshot: _AuditRecoverySnapshot | None = None + try: + current: PurgeAuditLog | None = session.get(PurgeAuditLog, record_id) + if current is None: + return "fallback" + snapshot = _capture_recovery_snapshot(current) + if current.status != STATUS_PENDING or current.trigger != TRIGGER_RETENTION: + return "retained" + predecessor: PurgeAuditLog | None = None + if current.entity_uuid is not None: + predecessor = _retention_predecessor(session, current) + suppressed: bool = _suppress_redundant_block(session, current, predecessor) + if not suppressed: + _retain_blocked(session, record_id) + session.commit() + return "suppressed" if suppressed else "retained" + except SQLAlchemyError: + session.rollback() + logger.warning( + "deletion_retention: persistence uncertainty finalizing blocked " + "audit row %s", + record_id, + exc_info=True, + ) + finally: + session.close() + return _recover_retention_blocked(record_id, snapshot) + + def _entity_exists(session: Session, record: PurgeAuditLog) -> bool | None: """Return whether the audit target exists, or None if it cannot resolve.""" # pylint: disable=import-outside-toplevel diff --git a/superset/migrations/versions/2026-08-06_18-00_b8d2f4a6c901_index_purge_audit_predecessor.py b/superset/migrations/versions/2026-08-06_18-00_b8d2f4a6c901_index_purge_audit_predecessor.py new file mode 100644 index 000000000000..7b7c34734b04 --- /dev/null +++ b/superset/migrations/versions/2026-08-06_18-00_b8d2f4a6c901_index_purge_audit_predecessor.py @@ -0,0 +1,43 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +"""Index purge audit predecessor lookups. + +Revision ID: b8d2f4a6c901 +Revises: d7cecc48bd55 +Create Date: 2026-08-06 18:00:00.000000 + +""" + +from superset.migrations.shared.utils import create_index, drop_index + +# revision identifiers, used by Alembic. +revision: str = "b8d2f4a6c901" +down_revision: str = "d7cecc48bd55" + +_INDEX_NAME: str = "ix_purge_audit_log_retention_predecessor" + + +def upgrade() -> None: + create_index( + "purge_audit_log", + _INDEX_NAME, + ["entity_uuid", "entity_type", "trigger", "created_on"], + ) + + +def downgrade() -> None: + drop_index("purge_audit_log", _INDEX_NAME) diff --git a/superset/models/purge_audit_log.py b/superset/models/purge_audit_log.py index 6bb3514af663..f4fe34de027a 100644 --- a/superset/models/purge_audit_log.py +++ b/superset/models/purge_audit_log.py @@ -52,6 +52,13 @@ class PurgeAuditLog(Model): # Backs reconcile_pending()'s stale-pending scan; mirrors the # index created by migration e7d93a524ff6. sa.Index("ix_purge_audit_log_status_created_on", "status", "created_on"), + sa.Index( + "ix_purge_audit_log_retention_predecessor", + "entity_uuid", + "entity_type", + "trigger", + "created_on", + ), ) id = Column(UUIDType(binary=True), primary_key=True, default=uuid4) diff --git a/superset/tasks/deletion_retention.py b/superset/tasks/deletion_retention.py index efb7119877d0..98431f0f17af 100644 --- a/superset/tasks/deletion_retention.py +++ b/superset/tasks/deletion_retention.py @@ -280,7 +280,17 @@ def _purge_one( removed_dashboard_slices=result.removed_dashboard_slices, ) elif result.blocked_reason is not None: - audit.block(record_id) + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(record_id) + ) + if disposition == "suppressed": + stats_logger_manager.instance.incr( + f"{_METRIC_PREFIX}.blocked_audit_suppressed" + ) + elif disposition == "fallback": + stats_logger_manager.instance.incr( + f"{_METRIC_PREFIX}.blocked_audit_dedupe_fallback" + ) else: audit.fail(record_id) return result diff --git a/tests/integration_tests/deletion_retention/audit_tests.py b/tests/integration_tests/deletion_retention/audit_tests.py index bf3c9946aeae..1394474e3253 100644 --- a/tests/integration_tests/deletion_retention/audit_tests.py +++ b/tests/integration_tests/deletion_retention/audit_tests.py @@ -19,9 +19,13 @@ from __future__ import annotations from datetime import datetime, timedelta -from unittest.mock import patch +from typing import Callable +from unittest.mock import MagicMock, patch from uuid import UUID +import pytest +from sqlalchemy.orm import Session + from superset import db from superset.commands.deletion_retention import audit from superset.commands.deletion_retention.audit import PurgeAuditLog @@ -32,6 +36,31 @@ class TestPurgeAudit(DeletionRetentionTestBase): + def _get_audit_record(self, record_id: UUID) -> PurgeAuditLog: + record: PurgeAuditLog | None = db.session.get(PurgeAuditLog, record_id) + assert record is not None + return record + + def _write_retention_record( + self, + *, + entity_uuid: str | None, + entity_type: str = "slices", + created_on: datetime | None = None, + ) -> UUID: + record_id: UUID | None = audit.write_ahead( + trigger=audit.TRIGGER_RETENTION, + actor=audit.ACTOR_SYSTEM, + entity_type=entity_type, + entity_uuid=entity_uuid, + ) + assert record_id is not None + if created_on is not None: + record: PurgeAuditLog = self._get_audit_record(record_id) + record.created_on = created_on + db.session.commit() + return record_id + def test_write_ahead_then_confirm(self) -> None: """A purge writes a pending audit row up front and flips it to confirmed after the delete commits.""" @@ -160,3 +189,304 @@ def test_blocked_attempt_does_not_keep_the_intended_removal_count(self) -> None: row = db.session.query(PurgeAuditLog).filter_by(id=record_id).one() assert row.status == audit.STATUS_BLOCKED assert row.removed_dashboard_slices == 0 + + def test_first_retention_block_is_retained(self) -> None: + record_id: UUID = self._write_retention_record(entity_uuid="first-block") + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(record_id) + ) + + record: PurgeAuditLog = self._get_audit_record(record_id) + assert disposition == "retained" + assert record.status == audit.STATUS_BLOCKED + + def test_repeated_retention_block_suppresses_current_provisional(self) -> None: + first_id: UUID = self._write_retention_record(entity_uuid="repeat-block") + audit.finalize_retention_blocked(first_id) + second_id: UUID = self._write_retention_record(entity_uuid="repeat-block") + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(second_id) + ) + + assert disposition == "suppressed" + first: PurgeAuditLog = self._get_audit_record(first_id) + assert first.status == audit.STATUS_BLOCKED + assert db.session.get(PurgeAuditLog, second_id) is None + + def test_equal_timestamp_is_ambiguous_and_retains_current(self) -> None: + timestamp: datetime = datetime.utcnow() + first_id: UUID = self._write_retention_record( + entity_uuid="equal-time", created_on=timestamp + ) + audit.finalize_retention_blocked(first_id) + second_id: UUID = self._write_retention_record( + entity_uuid="equal-time", created_on=timestamp + ) + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(second_id) + ) + + assert disposition == "retained" + second: PurgeAuditLog = self._get_audit_record(second_id) + assert second.status == audit.STATUS_BLOCKED + + def test_newer_concurrent_row_does_not_become_a_predecessor(self) -> None: + current_time: datetime = datetime.utcnow() + current_id: UUID = self._write_retention_record( + entity_uuid="overlap", created_on=current_time + ) + newer_id: UUID = self._write_retention_record( + entity_uuid="overlap", created_on=current_time + timedelta(seconds=1) + ) + audit.finalize_retention_blocked(newer_id) + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(current_id) + ) + + assert disposition == "retained" + current: PurgeAuditLog = self._get_audit_record(current_id) + assert current.status == audit.STATUS_BLOCKED + + def test_pending_predecessor_retains_current_block(self) -> None: + timestamp: datetime = datetime.utcnow() + self._write_retention_record(entity_uuid="pending-prior", created_on=timestamp) + current_id: UUID = self._write_retention_record( + entity_uuid="pending-prior", created_on=timestamp + timedelta(seconds=1) + ) + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(current_id) + ) + + assert disposition == "retained" + + def test_null_uuid_retains_blocked_evidence(self) -> None: + null_id: UUID = self._write_retention_record(entity_uuid=None) + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(null_id) + ) + + record: PurgeAuditLog = self._get_audit_record(null_id) + assert disposition == "retained" + assert record.status == audit.STATUS_BLOCKED + assert record.removed_dashboard_slices == 0 + + def test_same_uuid_across_entity_types_does_not_suppress(self) -> None: + chart_id: UUID = self._write_retention_record( + entity_uuid="shared-type", entity_type="slices" + ) + audit.finalize_retention_blocked(chart_id) + dashboard_id: UUID = self._write_retention_record( + entity_uuid="shared-type", entity_type="dashboards" + ) + + dashboard_disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(dashboard_id) + ) + + assert dashboard_disposition == "retained" + + def test_completed_current_record_is_immutable(self) -> None: + record_id: UUID = self._write_retention_record(entity_uuid="completed") + audit.fail(record_id) + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(record_id) + ) + + record: PurgeAuditLog = self._get_audit_record(record_id) + assert disposition == "retained" + assert record.status == audit.STATUS_FAILED + + def test_predecessor_lookup_failure_recovers_blocked_evidence(self) -> None: + record_id: UUID = self._write_retention_record(entity_uuid="lookup-failure") + + with patch( + "superset.commands.deletion_retention.audit._retention_predecessor", + side_effect=audit.SQLAlchemyError("lookup failed"), + ): + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(record_id) + ) + + record: PurgeAuditLog = self._get_audit_record(record_id) + assert disposition == "fallback" + assert record.status == audit.STATUS_BLOCKED + + def test_suppression_delete_failure_recovers_blocked_evidence(self) -> None: + first_id: UUID = self._write_retention_record(entity_uuid="delete-failure") + audit.finalize_retention_blocked(first_id) + current_id: UUID = self._write_retention_record(entity_uuid="delete-failure") + + with patch( + "superset.commands.deletion_retention.audit._suppress_redundant_block", + side_effect=audit.SQLAlchemyError("delete failed"), + ): + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(current_id) + ) + + record: PurgeAuditLog = self._get_audit_record(current_id) + assert disposition == "fallback" + assert record.status == audit.STATUS_BLOCKED + + def test_commit_failure_recovers_blocked_evidence(self) -> None: + record_id: UUID = self._write_retention_record(entity_uuid="commit-failure") + primary_session: Session = audit._dedicated_session() + recovery_session: Session = audit._dedicated_session() + with ( + patch.object( + primary_session, + "commit", + side_effect=audit.SQLAlchemyError("commit failed"), + ), + patch( + "superset.commands.deletion_retention.audit._dedicated_session", + side_effect=[primary_session, recovery_session], + ), + ): + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(record_id) + ) + + record: PurgeAuditLog = self._get_audit_record(record_id) + assert disposition == "fallback" + assert record.status == audit.STATUS_BLOCKED + + def test_uncertain_suppression_commit_recreates_absent_evidence(self) -> None: + first_id: UUID = self._write_retention_record(entity_uuid="absent-current") + audit.finalize_retention_blocked(first_id) + current_id: UUID = self._write_retention_record(entity_uuid="absent-current") + primary_session: Session = audit._dedicated_session() + recovery_session: Session = audit._dedicated_session() + primary_commit: Callable[[], None] = primary_session.commit + + def commit_then_raise() -> None: + primary_commit() + raise audit.SQLAlchemyError("commit acknowledgement lost") + + with ( + patch.object(primary_session, "commit", side_effect=commit_then_raise), + patch( + "superset.commands.deletion_retention.audit._dedicated_session", + side_effect=[primary_session, recovery_session], + ), + ): + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(current_id) + ) + + record: PurgeAuditLog = self._get_audit_record(current_id) + assert disposition == "fallback" + assert record.status == audit.STATUS_BLOCKED + + def test_failed_fallback_leaves_pending_evidence_for_reconciliation(self) -> None: + record_id: UUID = self._write_retention_record(entity_uuid="fallback-failure") + primary_session: Session = audit._dedicated_session() + recovery_session: Session = audit._dedicated_session() + with ( + patch.object( + primary_session, + "commit", + side_effect=audit.SQLAlchemyError("primary commit failed"), + ), + patch.object( + recovery_session, + "commit", + side_effect=audit.SQLAlchemyError("recovery commit failed"), + ), + patch( + "superset.commands.deletion_retention.audit._dedicated_session", + side_effect=[primary_session, recovery_session], + ), + ): + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(record_id) + ) + + record: PurgeAuditLog = self._get_audit_record(record_id) + assert disposition == "fallback" + assert record.status == audit.STATUS_PENDING + + def test_indeterminate_delete_rowcount_forces_persistence_recovery(self) -> None: + timestamp: datetime = datetime.utcnow() + current: PurgeAuditLog = PurgeAuditLog( + id=UUID("00000000-0000-0000-0000-000000000002"), + status=audit.STATUS_PENDING, + trigger=audit.TRIGGER_RETENTION, + actor=audit.ACTOR_SYSTEM, + entity_type="slices", + entity_uuid="indeterminate-rowcount", + created_on=timestamp, + ) + predecessor: PurgeAuditLog = PurgeAuditLog( + id=UUID("00000000-0000-0000-0000-000000000001"), + status=audit.STATUS_BLOCKED, + trigger=audit.TRIGGER_RETENTION, + actor=audit.ACTOR_SYSTEM, + entity_type="slices", + entity_uuid="indeterminate-rowcount", + created_on=timestamp - timedelta(seconds=1), + ) + result: MagicMock = MagicMock(rowcount=-1) + session: MagicMock = MagicMock() + session.execute.return_value = result + + with pytest.raises(audit.SQLAlchemyError, match="indeterminate"): + audit._suppress_redundant_block(session, current, predecessor) + + def test_overlap_duplicates_do_not_cause_unbounded_sequential_growth(self) -> None: + timestamp: datetime = datetime.utcnow() + first_id: UUID = self._write_retention_record( + entity_uuid="bounded-overlap", created_on=timestamp + ) + audit.finalize_retention_blocked(first_id) + overlap_id: UUID = self._write_retention_record( + entity_uuid="bounded-overlap", created_on=timestamp + ) + audit.finalize_retention_blocked(overlap_id) + later_id: UUID = self._write_retention_record( + entity_uuid="bounded-overlap", + created_on=timestamp + timedelta(seconds=1), + ) + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(later_id) + ) + + retained_count: int = ( + db.session.query(PurgeAuditLog) + .filter_by(entity_uuid="bounded-overlap") + .count() + ) + assert disposition == "suppressed" + assert retained_count == 2 + + def test_meaningful_outcome_transitions_start_new_blocked_periods(self) -> None: + transition: tuple[int, str] + for transition in enumerate( + ( + audit.STATUS_FAILED, + audit.STATUS_CONFIRMED, + audit.STATUS_TARGET_ABSENT, + ) + ): + index: int = transition[0] + status: str = transition[1] + entity_uuid: str = f"transition-{index}" + prior_id: UUID = self._write_retention_record(entity_uuid=entity_uuid) + audit.finalize(prior_id, status) + current_id: UUID = self._write_retention_record(entity_uuid=entity_uuid) + + disposition: audit.RetentionBlockedDisposition = ( + audit.finalize_retention_blocked(current_id) + ) + + current: PurgeAuditLog = self._get_audit_record(current_id) + assert disposition == "retained" + assert current.status == audit.STATUS_BLOCKED diff --git a/tests/integration_tests/deletion_retention/force_purge_tests.py b/tests/integration_tests/deletion_retention/force_purge_tests.py index 331b159f9ba6..f9907de0d04f 100644 --- a/tests/integration_tests/deletion_retention/force_purge_tests.py +++ b/tests/integration_tests/deletion_retention/force_purge_tests.py @@ -24,6 +24,7 @@ import pytest from superset import db +from superset.commands.deletion_retention import audit from superset.commands.deletion_retention.audit import PurgeAuditLog from superset.commands.deletion_retention.force_purge import ( AmbiguousPurgeTargetError, @@ -33,6 +34,7 @@ from superset.models.dashboard import Dashboard from superset.models.slice import Slice from superset.reports.models import ReportSchedule +from superset.tasks.deletion_retention import _purge_impl from ._base import DeletionRetentionTestBase @@ -145,6 +147,37 @@ def test_force_purge_preserves_report_reference_blocker(self) -> None: "associated alerts or reports exist", ) + def test_force_block_does_not_change_scheduled_deduplication_stream(self) -> None: + chart: Slice = self.make_chart("independent_block_streams") + report: ReportSchedule = ReportSchedule( + type="Report", + name="retention_it_independent_block_streams", + crontab="0 0 * * *", + chart=chart, + ) + db.session.add(report) + db.session.commit() + chart_uuid: str = str(chart.uuid) + self.soft_delete(chart, days_ago=90) + + first_scheduled: dict[str, object] = _purge_impl(30, dry_run=False) + force_result: dict[str, object] = ForcePurgeCommand(chart_uuid).run() + second_scheduled: dict[str, object] = _purge_impl(30, dry_run=False) + + records: list[PurgeAuditLog] = ( + db.session.query(PurgeAuditLog) + .filter_by(entity_uuid=chart_uuid) + .order_by(PurgeAuditLog.created_on) + .all() + ) + assert first_scheduled["blocked_by_reference"] == 1 + assert force_result["reason"] == "blocked" + assert second_scheduled["blocked_by_reference"] == 1 + assert [(record.trigger, record.status) for record in records] == [ + (audit.TRIGGER_RETENTION, audit.STATUS_BLOCKED), + (audit.TRIGGER_FORCE, audit.STATUS_BLOCKED), + ] + def test_force_purge_refuses_an_ambiguous_uuid(self) -> None: """A UUID matching two entity types is refused, not guessed. diff --git a/tests/integration_tests/deletion_retention/purge_tests.py b/tests/integration_tests/deletion_retention/purge_tests.py index bfca0868e2a3..ad8c3b1dbe50 100644 --- a/tests/integration_tests/deletion_retention/purge_tests.py +++ b/tests/integration_tests/deletion_retention/purge_tests.py @@ -47,6 +47,7 @@ from superset.models.user_attributes import UserAttribute from superset.reports.models import ReportSchedule from superset.tags.models import ObjectType, Tag, TaggedObject +from superset.tasks import deletion_retention as deletion_retention_task from superset.tasks.deletion_retention import _purge_impl from ._base import DeletionRetentionTestBase @@ -276,6 +277,80 @@ def test_report_reference_blocks_chart_purge(self) -> None: ) assert row.status == audit.STATUS_BLOCKED + def test_repeated_report_blocker_preserves_counts_and_suppresses_noise( + self, + ) -> None: + chart: Slice = self.make_chart("reported_repeatedly") + report: ReportSchedule = ReportSchedule( + type="Report", + name="retention_it_repeated_report", + crontab="0 0 * * *", + chart=chart, + ) + db.session.add(report) + db.session.commit() + chart_uuid: str = str(chart.uuid) + self.soft_delete(chart, days_ago=90) + + incr: MagicMock + gauge: MagicMock + with ( + patch.object( + deletion_retention_task.stats_logger_manager.instance, "incr" + ) as incr, + patch.object( + deletion_retention_task.stats_logger_manager.instance, "gauge" + ) as gauge, + ): + first_result: dict[str, Any] = _purge(window=30) + second_result: dict[str, Any] = _purge(window=30) + + assert first_result["blocked_by_reference"] == 1 + assert second_result["blocked_by_reference"] == 1 + assert ( + db.session.query(audit.PurgeAuditLog) + .filter_by( + entity_uuid=chart_uuid, + trigger=audit.TRIGGER_RETENTION, + status=audit.STATUS_BLOCKED, + ) + .count() + == 1 + ) + incr.assert_called_once_with("deletion_retention.blocked_audit_suppressed") + assert gauge.call_count == 2 + gauge.assert_called_with("deletion_retention.blocked_by_reference", 1) + + def test_blocked_audit_fallback_is_counted_without_changing_task_result( + self, + ) -> None: + chart: Slice = self.make_chart("reported_fallback") + report: ReportSchedule = ReportSchedule( + type="Report", + name="retention_it_fallback_report", + crontab="0 0 * * *", + chart=chart, + ) + db.session.add(report) + db.session.commit() + self.soft_delete(chart, days_ago=90) + + incr: MagicMock + with ( + patch.object( + audit, + "finalize_retention_blocked", + return_value="fallback", + ), + patch.object( + deletion_retention_task.stats_logger_manager.instance, "incr" + ) as incr, + ): + result: dict[str, Any] = _purge(window=30) + + assert result["blocked_by_reference"] == 1 + incr.assert_called_once_with("deletion_retention.blocked_audit_dedupe_fallback") + def test_restrictive_fk_blocks_dashboard_without_rewriting_referrer(self) -> None: """A welcome-dashboard FK remains authoritative during retention.""" dashboard = self.make_dashboard("welcome")