bito-code-review[bot] commented on code in PR #41549:
URL: https://github.com/apache/superset/pull/41549#discussion_r3658919532


##########
superset/tasks/deletion_retention.py:
##########
@@ -0,0 +1,258 @@
+# 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.
+"""Celery beat task: purge soft-deleted entities past the retention window.
+
+The deletion-domain analog of ``version_history.prune_old_versions``: where
+that ages out version rows while keeping the live entity, this removes
+entities that are already soft-deleted. For each
+``SoftDeleteMixin`` model it selects rows whose ``deleted_at`` is older than
+the per-workspace window and runs the shared cascade per entity, in bounded
+id-ordered batches. Convergent, not strictly idempotent: a re-run with the
+same clock and data removes nothing, but rows that have since crossed the
+cutoff are purged on a later run.
+"""
+
+from __future__ import annotations
+
+import logging
+from collections.abc import Iterator
+from datetime import datetime, timedelta
+from typing import Any, cast
+
+import sqlalchemy as sa
+from flask import current_app
+
+from superset import db
+from superset.commands.deletion_retention import audit
+from superset.commands.deletion_retention.purge_cascade import (
+    cascade_hard_delete,
+    CascadeResult,
+    dashboard_slice_count,
+    entity_uuid,
+    suppress_purge_association_versions,
+)
+from superset.commands.deletion_retention.window import 
resolve_retention_window
+from superset.extensions import celery_app, feature_flag_manager, 
stats_logger_manager
+from superset.models.helpers import (
+    skip_visibility_filter,
+    SoftDeleteMixin,
+)
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+_METRIC_PREFIX: str = "deletion_retention"
+# Below SQLite's historical 999 bind-variable limit; well under PostgreSQL and
+# MySQL limits. There is no existing SQLITE_MAX_VARIABLE_NUMBER symbol to 
reuse.
+_PURGE_DELETE_CHUNK: int = 500
+_BATCH: int = _PURGE_DELETE_CHUNK
+
+
+def _soft_delete_models() -> list[type[SoftDeleteMixin]]:
+    """The registered ``SoftDeleteMixin`` subclasses (dashboards, charts,
+    datasets), in a stable order."""
+    return list(SoftDeleteMixin._registered_subclasses)  # noqa: SLF001
+
+
+def _model_table(model: type[SoftDeleteMixin]) -> sa.Table:
+    """Return SQLAlchemy table metadata for a registered soft-delete model."""
+    return cast(sa.Table, cast(Any, model).__table__)
+
+
+def _model_table_name(model: type[SoftDeleteMixin]) -> str:
+    """Return the table name for a registered soft-delete model."""
+    return str(cast(Any, model).__tablename__)
+
+
+def _iter_eligible_ids(
+    model: type[SoftDeleteMixin], cutoff: datetime, batch: int
+) -> Iterator[list[int]]:
+    """Yield id-ordered batches of eligible row ids — ``deleted_at IS NOT NULL
+    AND deleted_at < cutoff`` — querying with the visibility-filter bypass so
+    soft-deleted rows are visible. Windowed by an ``id`` watermark so memory
+    and lock-hold stay bounded on a large first run."""
+    table = _model_table(model)
+    after_id = 0
+    while True:
+        with skip_visibility_filter(db.session, model):
+            ids = [
+                row[0]
+                for row in db.session.execute(
+                    sa.select(table.c.id)
+                    .where(table.c.deleted_at.is_not(None))
+                    .where(table.c.deleted_at < cutoff)
+                    .where(table.c.id > after_id)
+                    .order_by(table.c.id)
+                    .limit(batch)
+                )
+            ]
+        if not ids:
+            return
+        yield ids
+        if len(ids) < batch:
+            return
+        after_id = ids[-1]
+
+
+def _purge_impl(window_days: int, dry_run: bool) -> dict[str, Any]:
+    """Run one purge pass across all soft-delete models."""
+    if window_days <= 0:
+        logger.info("deletion_retention: window is 0 (disabled); skipping")
+        stats_logger_manager.instance.incr(f"{_METRIC_PREFIX}.skipped")
+        return {"skipped": 1}
+
+    cutoff = datetime.now() - timedelta(days=window_days)

Review Comment:
   <!-- Bito Reply -->
   The suggestion to use timezone-aware datetime is appropriate. Since the code 
now uses naive UTC to match existing `deleted_at` values, using 
`datetime.now(timezone.utc)` or `datetime.utcnow()` ensures consistency and 
avoids potential issues with local system time settings.
   
   **superset/tasks/deletion_retention.py**
   ```
   from datetime import datetime, timezone, timedelta
   
   # ...
   
       cutoff = datetime.now(timezone.utc) - timedelta(days=window_days)
   ```



##########
superset/commands/deletion_retention/audit.py:
##########
@@ -0,0 +1,225 @@
+# 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.
+"""Write-ahead purge audit record.
+
+Every purge — time-based or force — writes an immutable record that
+**survives** the entity it names, 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).
+
+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
+removed by the purge cascade.
+"""
+
+from __future__ import annotations
+
+import logging
+from datetime import datetime, timedelta
+from typing import Any, cast
+from uuid import UUID, uuid4
+
+import sqlalchemy as sa
+from sqlalchemy import Column, DateTime, Integer, String, Text
+from sqlalchemy.orm import Session, sessionmaker
+from sqlalchemy_utils import UUIDType
+
+from superset import db
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+
+def _dedicated_session() -> Session:
+    """A fresh session on its own connection, independent of the request /
+    task ``db.session``. The audit write must commit on its own so it survives
+    a rolled-back or crashed purge."""
+    return sessionmaker(bind=db.engine)()
+
+
+STATUS_PENDING = "pending"
+STATUS_CONFIRMED = "confirmed"
+STATUS_FAILED = "failed"
+STATUS_BLOCKED = "blocked"
+
+_PENDING_STALE_AFTER = timedelta(hours=1)
+
+TRIGGER_RETENTION = "retention"
+TRIGGER_FORCE = "force"
+
+ACTOR_SYSTEM = "system"
+
+
+class PurgeAuditLog(db.Model):
+    """Immutable, content-free record of a purge."""
+
+    __tablename__ = "purge_audit_log"
+
+    id = Column(UUIDType(binary=True), primary_key=True, default=uuid4)
+    status = Column(String(16), nullable=False, default=STATUS_PENDING)
+    trigger = Column(String(16), nullable=False)
+    actor = Column(String(256), nullable=False)
+    entity_type = Column(String(64), nullable=False)
+    entity_uuid = Column(String(36), nullable=True, index=True)
+    # Comma-joined UUIDs of charts left dangling / dashboards that lost a join
+    # row (force-purge visibility). Free text, content-free.
+    affected_referrers = Column(Text, nullable=True)
+    removed_dashboard_slices = Column(Integer, nullable=False, default=0)
+    created_on = Column(DateTime, nullable=False)
+    confirmed_on = Column(DateTime, nullable=True)
+
+
+def write_ahead(
+    *,
+    trigger: str,
+    actor: str,
+    entity_type: str,
+    entity_uuid: str | None,
+    removed_dashboard_slices: int = 0,
+) -> UUID | None:
+    """Insert a ``pending`` audit row on a dedicated session, before the
+    purge runs. Returns the row id to confirm later, or ``None`` if the audit
+    write itself fails (which must not block the purge)."""
+    session = _dedicated_session()
+    try:
+        record = PurgeAuditLog(
+            status=STATUS_PENDING,
+            trigger=trigger,
+            actor=actor,
+            entity_type=entity_type,
+            entity_uuid=entity_uuid,
+            removed_dashboard_slices=removed_dashboard_slices,
+            created_on=datetime.utcnow(),

Review Comment:
   <!-- Bito Reply -->
   The suggestion to replace `datetime.utcnow()` with 
`datetime.now(timezone.utc)` is correct. `datetime.utcnow()` is deprecated and 
returns a naive datetime object, whereas `datetime.now(timezone.utc)` provides 
a timezone-aware object, which is the recommended practice for handling UTC 
timestamps in modern Python.
   
   **superset/commands/deletion_retention/audit.py**
   ```
   created_on=datetime.now(timezone.utc),
   ```



##########
superset/commands/deletion_retention/audit.py:
##########
@@ -0,0 +1,225 @@
+# 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.
+"""Write-ahead purge audit record.
+
+Every purge — time-based or force — writes an immutable record that
+**survives** the entity it names, 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).
+
+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
+removed by the purge cascade.
+"""
+
+from __future__ import annotations
+
+import logging
+from datetime import datetime, timedelta
+from typing import Any, cast
+from uuid import UUID, uuid4
+
+import sqlalchemy as sa
+from sqlalchemy import Column, DateTime, Integer, String, Text
+from sqlalchemy.orm import Session, sessionmaker
+from sqlalchemy_utils import UUIDType
+
+from superset import db
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+
+def _dedicated_session() -> Session:
+    """A fresh session on its own connection, independent of the request /
+    task ``db.session``. The audit write must commit on its own so it survives
+    a rolled-back or crashed purge."""
+    return sessionmaker(bind=db.engine)()
+
+
+STATUS_PENDING = "pending"
+STATUS_CONFIRMED = "confirmed"
+STATUS_FAILED = "failed"
+STATUS_BLOCKED = "blocked"
+
+_PENDING_STALE_AFTER = timedelta(hours=1)
+
+TRIGGER_RETENTION = "retention"
+TRIGGER_FORCE = "force"
+
+ACTOR_SYSTEM = "system"
+
+
+class PurgeAuditLog(db.Model):
+    """Immutable, content-free record of a purge."""
+
+    __tablename__ = "purge_audit_log"
+
+    id = Column(UUIDType(binary=True), primary_key=True, default=uuid4)
+    status = Column(String(16), nullable=False, default=STATUS_PENDING)
+    trigger = Column(String(16), nullable=False)
+    actor = Column(String(256), nullable=False)
+    entity_type = Column(String(64), nullable=False)
+    entity_uuid = Column(String(36), nullable=True, index=True)
+    # Comma-joined UUIDs of charts left dangling / dashboards that lost a join
+    # row (force-purge visibility). Free text, content-free.
+    affected_referrers = Column(Text, nullable=True)
+    removed_dashboard_slices = Column(Integer, nullable=False, default=0)
+    created_on = Column(DateTime, nullable=False)
+    confirmed_on = Column(DateTime, nullable=True)
+
+
+def write_ahead(
+    *,
+    trigger: str,
+    actor: str,
+    entity_type: str,
+    entity_uuid: str | None,
+    removed_dashboard_slices: int = 0,
+) -> UUID | None:
+    """Insert a ``pending`` audit row on a dedicated session, before the
+    purge runs. Returns the row id to confirm later, or ``None`` if the audit
+    write itself fails (which must not block the purge)."""
+    session = _dedicated_session()
+    try:
+        record = PurgeAuditLog(
+            status=STATUS_PENDING,
+            trigger=trigger,
+            actor=actor,
+            entity_type=entity_type,
+            entity_uuid=entity_uuid,
+            removed_dashboard_slices=removed_dashboard_slices,
+            created_on=datetime.utcnow(),
+        )
+        session.add(record)
+        session.commit()
+        return cast(UUID, record.id)
+    except Exception:  # pylint: disable=broad-except
+        session.rollback()
+        logger.warning(
+            "deletion_retention: failed to write pending audit row", 
exc_info=True
+        )
+        return None
+    finally:
+        session.close()
+
+
+def finalize(record_id: UUID | None, status: str, **details: Any) -> None:
+    """Finalize a pending attempt on the dedicated audit session."""
+    if record_id is None:
+        return
+    session = _dedicated_session()
+    try:
+        record = session.get(PurgeAuditLog, record_id)
+        if record is None:
+            return
+        record.status = status
+        if status == STATUS_CONFIRMED:
+            record.confirmed_on = datetime.utcnow()

Review Comment:
   <!-- Bito Reply -->
   The suggestion to replace `datetime.utcnow()` with a timezone-aware 
alternative is valid. Using `datetime.utcnow()` is discouraged because it 
returns a naive datetime object, which can lead to ambiguity in 
timezone-sensitive operations. The recommended approach is to use 
`datetime.now(timezone.utc)` to ensure the datetime object is explicitly 
timezone-aware.
   
   **superset/commands/deletion_retention/audit.py**
   ```
   from datetime import datetime, timezone
   
   # ...
   
           record = PurgeAuditLog(
               # ...
               created_on=datetime.now(timezone.utc),
           )
   ```



##########
superset/commands/deletion_retention/audit.py:
##########
@@ -0,0 +1,225 @@
+# 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.
+"""Write-ahead purge audit record.
+
+Every purge — time-based or force — writes an immutable record that
+**survives** the entity it names, 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).
+
+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
+removed by the purge cascade.
+"""
+
+from __future__ import annotations
+
+import logging
+from datetime import datetime, timedelta
+from typing import Any, cast
+from uuid import UUID, uuid4
+
+import sqlalchemy as sa
+from sqlalchemy import Column, DateTime, Integer, String, Text
+from sqlalchemy.orm import Session, sessionmaker
+from sqlalchemy_utils import UUIDType
+
+from superset import db
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+
+def _dedicated_session() -> Session:
+    """A fresh session on its own connection, independent of the request /
+    task ``db.session``. The audit write must commit on its own so it survives
+    a rolled-back or crashed purge."""
+    return sessionmaker(bind=db.engine)()
+
+
+STATUS_PENDING = "pending"
+STATUS_CONFIRMED = "confirmed"
+STATUS_FAILED = "failed"
+STATUS_BLOCKED = "blocked"
+
+_PENDING_STALE_AFTER = timedelta(hours=1)
+
+TRIGGER_RETENTION = "retention"
+TRIGGER_FORCE = "force"
+
+ACTOR_SYSTEM = "system"
+
+
+class PurgeAuditLog(db.Model):
+    """Immutable, content-free record of a purge."""
+
+    __tablename__ = "purge_audit_log"
+
+    id = Column(UUIDType(binary=True), primary_key=True, default=uuid4)
+    status = Column(String(16), nullable=False, default=STATUS_PENDING)
+    trigger = Column(String(16), nullable=False)
+    actor = Column(String(256), nullable=False)
+    entity_type = Column(String(64), nullable=False)
+    entity_uuid = Column(String(36), nullable=True, index=True)
+    # Comma-joined UUIDs of charts left dangling / dashboards that lost a join
+    # row (force-purge visibility). Free text, content-free.
+    affected_referrers = Column(Text, nullable=True)
+    removed_dashboard_slices = Column(Integer, nullable=False, default=0)
+    created_on = Column(DateTime, nullable=False)
+    confirmed_on = Column(DateTime, nullable=True)
+
+
+def write_ahead(
+    *,
+    trigger: str,
+    actor: str,
+    entity_type: str,
+    entity_uuid: str | None,
+    removed_dashboard_slices: int = 0,
+) -> UUID | None:
+    """Insert a ``pending`` audit row on a dedicated session, before the
+    purge runs. Returns the row id to confirm later, or ``None`` if the audit
+    write itself fails (which must not block the purge)."""
+    session = _dedicated_session()
+    try:
+        record = PurgeAuditLog(
+            status=STATUS_PENDING,
+            trigger=trigger,
+            actor=actor,
+            entity_type=entity_type,
+            entity_uuid=entity_uuid,
+            removed_dashboard_slices=removed_dashboard_slices,
+            created_on=datetime.utcnow(),
+        )
+        session.add(record)
+        session.commit()
+        return cast(UUID, record.id)
+    except Exception:  # pylint: disable=broad-except
+        session.rollback()
+        logger.warning(
+            "deletion_retention: failed to write pending audit row", 
exc_info=True
+        )
+        return None
+    finally:
+        session.close()
+
+
+def finalize(record_id: UUID | None, status: str, **details: Any) -> None:
+    """Finalize a pending attempt on the dedicated audit session."""
+    if record_id is None:
+        return
+    session = _dedicated_session()
+    try:
+        record = session.get(PurgeAuditLog, record_id)
+        if record is None:
+            return
+        record.status = status
+        if status == STATUS_CONFIRMED:
+            record.confirmed_on = datetime.utcnow()
+        referrers = details.get("affected_referrers")
+        if referrers:
+            record.affected_referrers = ",".join(referrers)
+        removed_dashboard_slices = details.get("removed_dashboard_slices")
+        if removed_dashboard_slices is not None:
+            record.removed_dashboard_slices = removed_dashboard_slices
+        session.commit()
+    except Exception:  # pylint: disable=broad-except
+        session.rollback()
+        logger.warning(
+            "deletion_retention: failed to finalize audit row %s as %s",
+            record_id,
+            status,
+            exc_info=True,
+        )
+    finally:
+        session.close()
+
+
+def confirm(record_id: UUID | None, **details: Any) -> None:
+    """Mark an attempt confirmed after the entity transaction commits."""
+    finalize(record_id, STATUS_CONFIRMED, **details)
+
+
+def fail(record_id: UUID | None) -> None:
+    """Mark a known failed/no-op attempt so it does not remain pending."""
+    finalize(record_id, STATUS_FAILED)
+
+
+def block(record_id: UUID | None) -> None:
+    """Mark an attempt blocked by ordinary deletion policy."""
+    finalize(record_id, STATUS_BLOCKED)
+
+
+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
+    from superset.models.helpers import SoftDeleteMixin
+
+    if record.entity_uuid is None:
+        return None
+    for model in SoftDeleteMixin._registered_subclasses:  # noqa: SLF001
+        table = cast(Any, model).__table__
+        if table.name != record.entity_type or "uuid" not in table.c:
+            continue
+        return (
+            session.execute(
+                sa.select(table.c.id).where(table.c.uuid == 
record.entity_uuid).limit(1)
+            ).first()
+            is not None
+        )
+    return None
+
+
+def reconcile_pending(stale_before: datetime | None = None) -> dict[str, int]:
+    """Finalize stale pending attempts left by a process crash.
+
+    Missing entities prove the entity transaction committed, so the attempt is
+    confirmed. A surviving or unresolvable entity means the attempt did not
+    durably purge it and is finalized as failed; normal selection may retry.
+    """
+    cutoff = stale_before or datetime.utcnow() - _PENDING_STALE_AFTER
+    reconciled = confirmed = failed = 0
+    session = _dedicated_session()
+    try:
+        records = (
+            session.query(PurgeAuditLog)
+            .filter(PurgeAuditLog.status == STATUS_PENDING)
+            .filter(PurgeAuditLog.created_on < cutoff)
+            .all()
+        )
+        for record in records:
+            if _entity_exists(session, record) is False:
+                record.status = STATUS_CONFIRMED
+                record.confirmed_on = datetime.utcnow()

Review Comment:
   <!-- Bito Reply -->
   The suggestion to replace `datetime.utcnow()` with a timezone-aware 
alternative is correct. Using `datetime.utcnow()` is deprecated in many modern 
Python environments because it returns a naive datetime object, which can lead 
to ambiguity. Using a timezone-aware helper ensures consistency and avoids 
potential issues with time calculations across different environments.



##########
superset/commands/deletion_retention/audit.py:
##########
@@ -0,0 +1,231 @@
+# 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.
+"""Write-ahead purge audit record.
+
+Every purge — time-based or force — writes an immutable record that
+**survives** the entity it names, 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).
+
+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
+removed by the purge cascade.
+"""
+
+# Explicit commit/rollback on the dedicated session is the whole point of
+# this module — the audit row must survive independently of the purge
+# transaction, which the @transaction decorator (scoped to db.session)
+# cannot express.
+# pylint: disable=consider-using-transaction
+
+from __future__ import annotations
+
+import logging
+from datetime import datetime, timedelta
+from typing import Any, cast
+from uuid import UUID, uuid4
+
+import sqlalchemy as sa
+from sqlalchemy import Column, DateTime, Integer, String, Text
+from sqlalchemy.orm import Session, sessionmaker
+from sqlalchemy_utils import UUIDType
+
+from superset import db
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+
+def _dedicated_session() -> Session:
+    """A fresh session on its own connection, independent of the request /
+    task ``db.session``. The audit write must commit on its own so it survives
+    a rolled-back or crashed purge."""
+    return sessionmaker(bind=db.engine)()
+
+
+STATUS_PENDING = "pending"
+STATUS_CONFIRMED = "confirmed"
+STATUS_FAILED = "failed"
+STATUS_BLOCKED = "blocked"
+
+_PENDING_STALE_AFTER = timedelta(hours=1)
+
+TRIGGER_RETENTION = "retention"
+TRIGGER_FORCE = "force"
+
+ACTOR_SYSTEM = "system"
+
+
+class PurgeAuditLog(db.Model):
+    """Immutable, content-free record of a purge."""
+
+    __tablename__ = "purge_audit_log"
+
+    id = Column(UUIDType(binary=True), primary_key=True, default=uuid4)
+    status = Column(String(16), nullable=False, default=STATUS_PENDING)
+    trigger = Column(String(16), nullable=False)
+    actor = Column(String(256), nullable=False)
+    entity_type = Column(String(64), nullable=False)
+    entity_uuid = Column(String(36), nullable=True, index=True)
+    # Comma-joined UUIDs of charts left dangling / dashboards that lost a join
+    # row (force-purge visibility). Free text, content-free.
+    affected_referrers = Column(Text, nullable=True)
+    removed_dashboard_slices = Column(Integer, nullable=False, default=0)
+    created_on = Column(DateTime, nullable=False)
+    confirmed_on = Column(DateTime, nullable=True)
+
+
+def write_ahead(
+    *,
+    trigger: str,
+    actor: str,
+    entity_type: str,
+    entity_uuid: str | None,
+    removed_dashboard_slices: int = 0,
+) -> UUID | None:
+    """Insert a ``pending`` audit row on a dedicated session, before the
+    purge runs. Returns the row id to confirm later, or ``None`` if the audit
+    write itself fails (which must not block the purge)."""
+    session = _dedicated_session()
+    try:
+        record = PurgeAuditLog(
+            status=STATUS_PENDING,
+            trigger=trigger,
+            actor=actor,
+            entity_type=entity_type,
+            entity_uuid=entity_uuid,
+            removed_dashboard_slices=removed_dashboard_slices,
+            created_on=datetime.utcnow(),
+        )
+        session.add(record)
+        session.commit()
+        return cast(UUID, record.id)
+    except Exception:  # pylint: disable=broad-except
+        session.rollback()
+        logger.warning(
+            "deletion_retention: failed to write pending audit row", 
exc_info=True
+        )
+        return None
+    finally:
+        session.close()
+
+
+def finalize(record_id: UUID | None, status: str, **details: Any) -> None:
+    """Finalize a pending attempt on the dedicated audit session."""
+    if record_id is None:
+        return
+    session = _dedicated_session()
+    try:
+        record = session.get(PurgeAuditLog, record_id)
+        if record is None:
+            return
+        record.status = status
+        if status == STATUS_CONFIRMED:
+            record.confirmed_on = datetime.utcnow()
+        referrers = details.get("affected_referrers")
+        if referrers:
+            record.affected_referrers = ",".join(referrers)

Review Comment:
   <!-- Bito Reply -->
   The suggestion to replace `datetime.utcnow()` with 
`datetime.now(timezone.utc)` is correct and improves the code by using a 
timezone-aware datetime object, which is a best practice in Python. Since you 
have already implemented this fix using the shared `_utc_now()` helper, the 
suggestion has been addressed.
   
   **superset/commands/deletion_retention/audit.py**
   ```
   record.status = status
           if status == STATUS_CONFIRMED:
               record.confirmed_on = datetime.now(timezone.utc)
           referrers = details.get("affected_referrers")
           if referrers:
               record.affected_referrers = ",".join(referrers)
   ```



##########
superset/commands/deletion_retention/audit.py:
##########
@@ -0,0 +1,231 @@
+# 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.
+"""Write-ahead purge audit record.
+
+Every purge — time-based or force — writes an immutable record that
+**survives** the entity it names, 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).
+
+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
+removed by the purge cascade.
+"""
+
+# Explicit commit/rollback on the dedicated session is the whole point of
+# this module — the audit row must survive independently of the purge
+# transaction, which the @transaction decorator (scoped to db.session)
+# cannot express.
+# pylint: disable=consider-using-transaction
+
+from __future__ import annotations
+
+import logging
+from datetime import datetime, timedelta
+from typing import Any, cast
+from uuid import UUID, uuid4
+
+import sqlalchemy as sa
+from sqlalchemy import Column, DateTime, Integer, String, Text
+from sqlalchemy.orm import Session, sessionmaker
+from sqlalchemy_utils import UUIDType
+
+from superset import db
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+
+def _dedicated_session() -> Session:
+    """A fresh session on its own connection, independent of the request /
+    task ``db.session``. The audit write must commit on its own so it survives
+    a rolled-back or crashed purge."""
+    return sessionmaker(bind=db.engine)()
+
+
+STATUS_PENDING = "pending"
+STATUS_CONFIRMED = "confirmed"
+STATUS_FAILED = "failed"
+STATUS_BLOCKED = "blocked"
+
+_PENDING_STALE_AFTER = timedelta(hours=1)
+
+TRIGGER_RETENTION = "retention"
+TRIGGER_FORCE = "force"
+
+ACTOR_SYSTEM = "system"
+
+
+class PurgeAuditLog(db.Model):
+    """Immutable, content-free record of a purge."""
+
+    __tablename__ = "purge_audit_log"
+
+    id = Column(UUIDType(binary=True), primary_key=True, default=uuid4)
+    status = Column(String(16), nullable=False, default=STATUS_PENDING)
+    trigger = Column(String(16), nullable=False)
+    actor = Column(String(256), nullable=False)
+    entity_type = Column(String(64), nullable=False)
+    entity_uuid = Column(String(36), nullable=True, index=True)
+    # Comma-joined UUIDs of charts left dangling / dashboards that lost a join
+    # row (force-purge visibility). Free text, content-free.
+    affected_referrers = Column(Text, nullable=True)
+    removed_dashboard_slices = Column(Integer, nullable=False, default=0)
+    created_on = Column(DateTime, nullable=False)
+    confirmed_on = Column(DateTime, nullable=True)
+
+
+def write_ahead(
+    *,
+    trigger: str,
+    actor: str,
+    entity_type: str,
+    entity_uuid: str | None,
+    removed_dashboard_slices: int = 0,
+) -> UUID | None:
+    """Insert a ``pending`` audit row on a dedicated session, before the
+    purge runs. Returns the row id to confirm later, or ``None`` if the audit
+    write itself fails (which must not block the purge)."""
+    session = _dedicated_session()
+    try:
+        record = PurgeAuditLog(
+            status=STATUS_PENDING,
+            trigger=trigger,
+            actor=actor,
+            entity_type=entity_type,
+            entity_uuid=entity_uuid,
+            removed_dashboard_slices=removed_dashboard_slices,
+            created_on=datetime.utcnow(),
+        )
+        session.add(record)
+        session.commit()
+        return cast(UUID, record.id)
+    except Exception:  # pylint: disable=broad-except
+        session.rollback()
+        logger.warning(
+            "deletion_retention: failed to write pending audit row", 
exc_info=True
+        )
+        return None
+    finally:
+        session.close()
+
+
+def finalize(record_id: UUID | None, status: str, **details: Any) -> None:
+    """Finalize a pending attempt on the dedicated audit session."""
+    if record_id is None:
+        return
+    session = _dedicated_session()
+    try:
+        record = session.get(PurgeAuditLog, record_id)
+        if record is None:
+            return
+        record.status = status
+        if status == STATUS_CONFIRMED:
+            record.confirmed_on = datetime.utcnow()
+        referrers = details.get("affected_referrers")
+        if referrers:
+            record.affected_referrers = ",".join(referrers)
+        removed_dashboard_slices = details.get("removed_dashboard_slices")
+        if removed_dashboard_slices is not None:
+            record.removed_dashboard_slices = removed_dashboard_slices
+        session.commit()
+    except Exception:  # pylint: disable=broad-except
+        session.rollback()
+        logger.warning(
+            "deletion_retention: failed to finalize audit row %s as %s",
+            record_id,
+            status,
+            exc_info=True,
+        )
+    finally:
+        session.close()
+
+
+def confirm(record_id: UUID | None, **details: Any) -> None:
+    """Mark an attempt confirmed after the entity transaction commits."""
+    finalize(record_id, STATUS_CONFIRMED, **details)
+
+
+def fail(record_id: UUID | None) -> None:
+    """Mark a known failed/no-op attempt so it does not remain pending."""
+    finalize(record_id, STATUS_FAILED)
+
+
+def block(record_id: UUID | None) -> None:
+    """Mark an attempt blocked by ordinary deletion policy."""
+    finalize(record_id, STATUS_BLOCKED)
+
+
+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
+    from superset.models.helpers import SoftDeleteMixin
+
+    if record.entity_uuid is None:
+        return None
+    for model in SoftDeleteMixin._registered_subclasses:  # noqa: SLF001
+        table = cast(Any, model).__table__
+        if table.name != record.entity_type or "uuid" not in table.c:
+            continue
+        return (
+            session.execute(
+                sa.select(table.c.id).where(table.c.uuid == 
record.entity_uuid).limit(1)
+            ).first()
+            is not None
+        )
+    return None
+
+
+def reconcile_pending(stale_before: datetime | None = None) -> dict[str, int]:
+    """Finalize stale pending attempts left by a process crash.
+
+    Missing entities prove the entity transaction committed, so the attempt is
+    confirmed. A surviving or unresolvable entity means the attempt did not
+    durably purge it and is finalized as failed; normal selection may retry.
+    """
+    cutoff = stale_before or datetime.utcnow() - _PENDING_STALE_AFTER
+    reconciled = confirmed = failed = 0
+    session = _dedicated_session()
+    try:
+        records = (
+            session.query(PurgeAuditLog)
+            .filter(PurgeAuditLog.status == STATUS_PENDING)
+            .filter(PurgeAuditLog.created_on < cutoff)
+            .all()
+        )
+        for record in records:
+            if _entity_exists(session, record) is False:
+                record.status = STATUS_CONFIRMED
+                record.confirmed_on = datetime.utcnow()
+                confirmed += 1
+            else:
+                record.status = STATUS_FAILED

Review Comment:
   <!-- Bito Reply -->
   The suggestion provided by the reviewer is appropriate and improves the code 
by using a more robust and standard way to handle UTC timestamps in Python. 
Applying this change ensures consistency with the shared helper and avoids 
potential issues associated with the deprecated `datetime.utcnow()` method.
   
   **superset/commands/deletion_retention/audit.py**
   ```
   for record in records:
               if _entity_exists(session, record) is False:
                   record.status = STATUS_CONFIRMED
                   record.confirmed_on = datetime.now(timezone.utc)
                   confirmed += 1
               else:
                   record.status = STATUS_FAILED
   ```



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