sadpandajoe commented on code in PR #44892:
URL: https://github.com/apache/superset/pull/44892#discussion_r4183962399


##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -860,16 +894,131 @@ def declare(
             cleanup_permission=cleanup_dataset_permission,
         ),
     }
-    return validate_unique_root_policies(registry.values())
+    return tuple(registry.values())
+
+
+def _host_purge_policies(
+    provider: Callable[[], Any] | None,
+) -> tuple[PurgeEntityPolicy, ...]:
+    """Return the purge policies a host installed for its own roots.
+
+    A host distribution can carry ``SoftDeleteMixin`` entities this package
+    cannot import. The retention task discovers those entities through the
+    mixin registry, so without a policy they reach the cascade as an
+    unsupported model. The host therefore declares their purge behavior and
+    installs it under ``PURGE_POLICIES_FUNC``.
+
+    Host boundary: everything about a host policy is settled here, where it is
+    admitted -- an unavailable provider, a malformed payload, a collision with
+    a root declared in this package, two declarations for one root, and a
+    declaration the shared cleanup could not execute. Each is logged and
+    dropped, so that root is reported as unsupported instead of failing row by
+    row in a scheduled run. A broken host declaration must not take the purge
+    down with it, and must never redefine how a chart, dashboard or dataset is
+    purged.
+    """
+    if provider is None:
+        return ()
+    try:
+        provided: Any = provider()
+    except Exception:  # pylint: disable=broad-except
+        logger.exception(
+            "purge_policy: %s is unavailable; keeping built-in roots only",
+            HOST_POLICIES_CONFIG_KEY,
+        )
+        return ()
+    if not isinstance(provided, (list, tuple)) or not all(
+        isinstance(policy, PurgeEntityPolicy) for policy in provided
+    ):
+        logger.error(
+            "purge_policy: %s returned %s; expected a sequence of 
PurgeEntityPolicy",
+            HOST_POLICIES_CONFIG_KEY,
+            type(provided).__name__,
+        )
+        return ()
+    builtin_roots: frozenset[type[Any]] = frozenset(
+        policy.model for policy in _builtin_purge_policies()
+    )
+    declared: dict[type[Any], int] = {}
+    for policy in provided:
+        declared[policy.model] = declared.get(policy.model, 0) + 1
+    # Which of two declarations for one root is authoritative is undecidable,
+    # and the loser would still delete rows. Dropping both leaves the model
+    # reported as unsupported, which is the recoverable outcome. Resolving it
+    # here also keeps the duplicate away from validate_unique_root_policies,
+    # whose ValueError would abort the whole scheduled run.
+    duplicated: list[type[Any]] = [
+        model for model, count in declared.items() if count > 1
+    ]
+    for model in duplicated:
+        logger.error(
+            "purge_policy: %s declares %d policies for %s; ignoring all of 
them",
+            HOST_POLICIES_CONFIG_KEY,
+            declared[model],
+            model.__name__,
+        )
+    accepted: list[PurgeEntityPolicy] = []
+    for policy in provided:
+        if policy.model in builtin_roots:
+            logger.error(
+                "purge_policy: host policy for built-in root %s ignored",
+                policy.model.__name__,
+            )
+            continue
+        if policy.model in duplicated:
+            continue
+        try:
+            _validated_policy(policy)
+        except Exception as ex:  # pylint: disable=broad-except
+            # Validated at admission rather than when a row is purged: a
+            # declaration the cleanup cannot execute would otherwise fail once
+            # per eligible row, counted as a cascade failure, and stay
+            # invisible until something aged past the window.
+            logger.error(
+                "purge_policy: host policy for %s rejected: %s",
+                policy.model.__name__,
+                ex,
+            )
+            continue
+        accepted.append(policy)
+    return tuple(accepted)
+
+
+def purge_policy_registry() -> Mapping[type[Any], PurgeEntityPolicy]:
+    """Index the built-in purge roots plus any the host installed.
+
+    The index is cached per installed provider, not once per process. Freezing
+    it at first use would pin whatever happened to be installed at that
+    moment -- including nothing at all, for a call made before startup
+    finished -- until a restart. Keying the cache on the provider means it is
+    invoked once however many roots are purged, a host that installs or
+    replaces one is honored, and the ordinary case (no provider) resolves to a
+    single cached index.
+    """
+    provider: Callable[[], Any] | None = (
+        current_app.config.get(HOST_POLICIES_CONFIG_KEY) if has_app_context() 
else None
+    )
+    return _registry_for(provider)

Review Comment:
   A callable provider implemented as a mutable dataclass raises `TypeError: 
unhashable type` at this cache before it can be invoked or isolated, so even 
the built-in roots stop purging. Could this cache by provider identity without 
requiring extension callables to be hashable?



##########
superset/tasks/deletion_retention.py:
##########
@@ -132,6 +133,54 @@ def _report_model_counts(outcome: str, counts: dict[str, 
int]) -> None:
         )
 
 
+@dataclass
+class _PassTotals:
+    """What one pass over the soft-delete roots produced."""
+
+    purged: dict[str, int] = field(default_factory=dict)
+    would_purge: dict[str, int] = field(default_factory=dict)
+    unsupported: dict[str, int] = field(default_factory=dict)
+    cascade_failures: int = 0
+    blocked: int = 0
+    scan_failures: int = 0
+
+
+def _purge_roots(cutoff: datetime, dry_run: bool) -> _PassTotals:
+    """Process each registered root, isolating one root's failure from the 
rest."""
+    totals = _PassTotals()
+    for model in _soft_delete_models():
+        entity_type = _model_table_name(model)
+        try:
+            if model not in purge_policy_registry():
+                totals.unsupported[entity_type] = 1
+                logger.warning(
+                    "deletion_retention: skipping %s: no purge policy", 
entity_type
+                )
+                stats_logger_manager.instance.incr(
+                    f"{_METRIC_PREFIX}.unsupported_models.{entity_type}"
+                )
+                continue
+            purged_n, would_n, failed_n, blocked_n = _purge_model(
+                model, cutoff, dry_run
+            )
+        except Exception:  # pylint: disable=broad-except
+            # One root must not cost the others their run. The eligible-id scan
+            # runs outside _purge_model's per-entity handler, so a column it
+            # cannot read or a transient database error would otherwise abort
+            # the whole pass -- including the roots that would have purged.
+            db.session.rollback()  # pylint: disable=consider-using-transaction

Review Comment:
   If the first scan page purges 500 rows and the next SELECT fails, those 
deletions are already committed, but this handler discards the root’s 
accumulated counts; the summary and purge gauges report only later roots. Could 
scan isolation preserve the counts accumulated before the failing page?



##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -905,32 +1056,226 @@ def _validated_purge_policy(model: type[Any]) -> 
PurgeEntityPolicy:
             details = (
                 f"{details}; " if details else ""
             ) + f"stale_listeners=[{', '.join(coverage.stale_listeners)}]"
-        raise RuntimeError(f"Incomplete purge policy for {model.__name__}: 
{details}")
+        raise RuntimeError(
+            f"Incomplete purge policy for {policy.model.__name__}: {details}"
+        )
     return policy
 
 
+def _ownership_dependency(
+    policy: PurgeEntityPolicy, related_table: str
+) -> DependencyPolicy | None:
+    """The single owned/association edge attaching *related_table*, if 
clear."""
+    candidates: tuple[DependencyPolicy, ...] = tuple(
+        dependency
+        for dependency in policy.dependencies
+        if dependency.classification
+        in {DependencyClassification.OWNED, 
DependencyClassification.ASSOCIATION}
+        and dependency.key.related_table == related_table
+        and dependency.key.direction == "inbound"
+    )
+    return candidates[0] if len(candidates) == 1 else None
+
+
+def _validate_owned_traversal(policy: PurgeEntityPolicy) -> None:
+    """Reject an owned table reachable only through an association.
+
+    ``cascade_hard_delete`` empties associations before owned children, while
+    an owned table's predicate selects its rows *through* its ownership path
+    (see ``_owner_value_select``). If a hop on that path is an association,
+    its rows are already gone when the owned delete runs: the statement
+    matches nothing, leaving the descendants orphaned where foreign keys are
+    unenforced and blocking the root's delete where they are not.
+
+    Refused at declaration time rather than executed. The shape has a
+    remedy -- classify the intermediate table as owned, which places it in
+    the same phase as what it leads to.
+    """
+    root_table: str = sa.inspect(policy.model).local_table.name
+    for dependency in policy.dependencies:
+        if dependency.classification is not DependencyClassification.OWNED:
+            continue
+        table_name: str = dependency.key.owner_table
+        visited: set[str] = set()
+        while table_name != root_table and table_name not in visited:
+            visited.add(table_name)
+            hop: DependencyPolicy | None = _ownership_dependency(policy, 
table_name)
+            if hop is None:
+                # An absent or ambiguous path is reported by coverage and by
+                # _ownership_edge at execution; not this check's business.
+                break
+            if hop.classification is DependencyClassification.ASSOCIATION:
+                raise RuntimeError(
+                    f"Owned dependency {dependency.key.describe()} is 
reachable "
+                    f"only through association {hop.key.describe()}; "
+                    "associations are deleted first, so the owned rows would "
+                    "be orphaned"
+                )
+            table_name = hop.key.owner_table
+
+
+def _tables_share_a_foreign_key(metadata: sa.MetaData, first: str, second: 
str) -> bool:
+    """Whether either table declares a foreign key into the other."""
+    for owner_name, other_name in ((first, second), (second, first)):
+        table: sa.Table | None = metadata.tables.get(owner_name)
+        if table is None:
+            continue
+        for constraint in table.foreign_key_constraints:
+            if constraint.elements[0].column.table.name == other_name:
+                return True
+    return False
+
+
+def _validate_scanner_requirements(policy: PurgeEntityPolicy) -> None:
+    """Reject a root the scheduled scan cannot page through.
+
+    The retention task selects, windows and orders eligible rows by ``id``,
+    and the cascade pins each row by it. A root keyed on something else -- a
+    UUID primary key with no ``id`` column -- raises inside the scan, outside
+    the per-row error handling, so the run aborts before the remaining roots
+    are reached.
+    """
+    table: sa.Table = sa.inspect(policy.model).local_table
+    if "id" not in table.c:
+        raise RuntimeError(
+            f"Purge root {policy.model.__name__} has no 'id' column; the "
+            "scheduled scan pages eligible rows by id"
+        )
+    if "deleted_at" not in table.c:
+        raise RuntimeError(
+            f"Purge root {policy.model.__name__} has no 'deleted_at' column; "
+            "every purge path requires the row to be archived first"
+        )
+
+
+def _validate_recursive_ownership(policy: PurgeEntityPolicy) -> None:
+    """Reject a self-referencing owned table under the stock cleanup.
+
+    ``delete_owned_children`` issues one statement per declared edge, so a
+    table that owns itself is pruned one level deep: a three-level tree
+    leaves its deepest rows behind -- orphaned where foreign keys are
+    unenforced, and blocking the root's own delete where they are not. A host
+    declaring this shape supplies cleanup that walks the whole sub-tree.
+    """
+    for dependency in policy.dependencies:
+        if dependency.classification is not DependencyClassification.OWNED:
+            continue
+        key: DependencyKey = dependency.key
+        if key.owner_table != key.related_table:
+            continue
+        if policy.delete_owned_children is delete_owned_children:
+            raise RuntimeError(
+                f"Owned dependency {key.describe()} is self-referencing; the "
+                "stock owned-child cleanup deletes one level, so this policy "
+                "must supply its own delete_owned_children"
+            )
+
+
+def _validate_sibling_ownership(policy: PurgeEntityPolicy) -> None:
+    """Reject owned siblings that reference one another.
+
+    Owned deletes are sorted by ownership depth with a stable sort, so two
+    owned tables at the same depth are deleted in the order they were
+    declared. With a foreign key between them, that order decides the
+    outcome: removing the referenced table first is refused by an enforced
+    constraint and rolls the purge back. Ordering deletes by their own
+    foreign keys would be the richer fix; refusing the shape keeps the
+    engine's contract honest until something needs it.
+    """
+    metadata: sa.MetaData = sa.inspect(policy.model).local_table.metadata
+    by_depth: dict[int, list[DependencyPolicy]] = {}
+    for dependency in policy.dependencies:
+        if dependency.classification is not DependencyClassification.OWNED:
+            continue
+        depth: int = _dependency_owner_depth(policy, dependency.key)

Review Comment:
   For `workflow → OWNED node` plus `node → OWNED node`, a custom recursive 
cleanup passes the self-reference guard, but this stock-depth calculation then 
raises `Expected one ownership path to node, found 2` and drops the policy as 
unsupported. Could custom-cleaned child trees avoid this stock traversal 
requirement, as self-referencing roots already do?



##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -905,32 +1056,226 @@ def _validated_purge_policy(model: type[Any]) -> 
PurgeEntityPolicy:
             details = (
                 f"{details}; " if details else ""
             ) + f"stale_listeners=[{', '.join(coverage.stale_listeners)}]"
-        raise RuntimeError(f"Incomplete purge policy for {model.__name__}: 
{details}")
+        raise RuntimeError(
+            f"Incomplete purge policy for {policy.model.__name__}: {details}"
+        )
     return policy
 
 
+def _ownership_dependency(
+    policy: PurgeEntityPolicy, related_table: str
+) -> DependencyPolicy | None:
+    """The single owned/association edge attaching *related_table*, if 
clear."""
+    candidates: tuple[DependencyPolicy, ...] = tuple(
+        dependency
+        for dependency in policy.dependencies
+        if dependency.classification
+        in {DependencyClassification.OWNED, 
DependencyClassification.ASSOCIATION}
+        and dependency.key.related_table == related_table
+        and dependency.key.direction == "inbound"
+    )
+    return candidates[0] if len(candidates) == 1 else None
+
+
+def _validate_owned_traversal(policy: PurgeEntityPolicy) -> None:
+    """Reject an owned table reachable only through an association.
+
+    ``cascade_hard_delete`` empties associations before owned children, while
+    an owned table's predicate selects its rows *through* its ownership path
+    (see ``_owner_value_select``). If a hop on that path is an association,
+    its rows are already gone when the owned delete runs: the statement
+    matches nothing, leaving the descendants orphaned where foreign keys are
+    unenforced and blocking the root's delete where they are not.
+
+    Refused at declaration time rather than executed. The shape has a
+    remedy -- classify the intermediate table as owned, which places it in
+    the same phase as what it leads to.
+    """
+    root_table: str = sa.inspect(policy.model).local_table.name
+    for dependency in policy.dependencies:
+        if dependency.classification is not DependencyClassification.OWNED:
+            continue
+        table_name: str = dependency.key.owner_table
+        visited: set[str] = set()
+        while table_name != root_table and table_name not in visited:
+            visited.add(table_name)
+            hop: DependencyPolicy | None = _ownership_dependency(policy, 
table_name)
+            if hop is None:
+                # An absent or ambiguous path is reported by coverage and by
+                # _ownership_edge at execution; not this check's business.
+                break
+            if hop.classification is DependencyClassification.ASSOCIATION:
+                raise RuntimeError(
+                    f"Owned dependency {dependency.key.describe()} is 
reachable "
+                    f"only through association {hop.key.describe()}; "
+                    "associations are deleted first, so the owned rows would "
+                    "be orphaned"
+                )
+            table_name = hop.key.owner_table
+
+
+def _tables_share_a_foreign_key(metadata: sa.MetaData, first: str, second: 
str) -> bool:
+    """Whether either table declares a foreign key into the other."""
+    for owner_name, other_name in ((first, second), (second, first)):
+        table: sa.Table | None = metadata.tables.get(owner_name)
+        if table is None:
+            continue
+        for constraint in table.foreign_key_constraints:
+            if constraint.elements[0].column.table.name == other_name:
+                return True
+    return False
+
+
+def _validate_scanner_requirements(policy: PurgeEntityPolicy) -> None:
+    """Reject a root the scheduled scan cannot page through.
+
+    The retention task selects, windows and orders eligible rows by ``id``,
+    and the cascade pins each row by it. A root keyed on something else -- a
+    UUID primary key with no ``id`` column -- raises inside the scan, outside
+    the per-row error handling, so the run aborts before the remaining roots
+    are reached.
+    """
+    table: sa.Table = sa.inspect(policy.model).local_table
+    if "id" not in table.c:
+        raise RuntimeError(
+            f"Purge root {policy.model.__name__} has no 'id' column; the "
+            "scheduled scan pages eligible rows by id"
+        )
+    if "deleted_at" not in table.c:
+        raise RuntimeError(
+            f"Purge root {policy.model.__name__} has no 'deleted_at' column; "
+            "every purge path requires the row to be archived first"
+        )
+
+
+def _validate_recursive_ownership(policy: PurgeEntityPolicy) -> None:
+    """Reject a self-referencing owned table under the stock cleanup.
+
+    ``delete_owned_children`` issues one statement per declared edge, so a
+    table that owns itself is pruned one level deep: a three-level tree
+    leaves its deepest rows behind -- orphaned where foreign keys are
+    unenforced, and blocking the root's own delete where they are not. A host
+    declaring this shape supplies cleanup that walks the whole sub-tree.
+    """
+    for dependency in policy.dependencies:
+        if dependency.classification is not DependencyClassification.OWNED:
+            continue
+        key: DependencyKey = dependency.key
+        if key.owner_table != key.related_table:
+            continue
+        if policy.delete_owned_children is delete_owned_children:
+            raise RuntimeError(
+                f"Owned dependency {key.describe()} is self-referencing; the "
+                "stock owned-child cleanup deletes one level, so this policy "
+                "must supply its own delete_owned_children"
+            )
+
+
+def _validate_sibling_ownership(policy: PurgeEntityPolicy) -> None:
+    """Reject owned siblings that reference one another.
+
+    Owned deletes are sorted by ownership depth with a stable sort, so two
+    owned tables at the same depth are deleted in the order they were
+    declared. With a foreign key between them, that order decides the
+    outcome: removing the referenced table first is refused by an enforced
+    constraint and rolls the purge back. Ordering deletes by their own
+    foreign keys would be the richer fix; refusing the shape keeps the
+    engine's contract honest until something needs it.
+    """
+    metadata: sa.MetaData = sa.inspect(policy.model).local_table.metadata
+    by_depth: dict[int, list[DependencyPolicy]] = {}
+    for dependency in policy.dependencies:
+        if dependency.classification is not DependencyClassification.OWNED:
+            continue
+        depth: int = _dependency_owner_depth(policy, dependency.key)
+        by_depth.setdefault(depth, []).append(dependency)
+    for siblings in by_depth.values():
+        for index, first in enumerate(siblings):
+            for second in siblings[index + 1 :]:
+                if first.key.related_table == second.key.related_table:
+                    continue
+                if _tables_share_a_foreign_key(
+                    metadata, first.key.related_table, second.key.related_table
+                ):
+                    raise RuntimeError(
+                        f"Owned dependencies {first.key.describe()} and "
+                        f"{second.key.describe()} reference each other; owned "
+                        "deletes run in declaration order, so whether this "
+                        "graph purges depends on it"
+                    )
+
+
+#: Classifications the shared cleanup executes as SQL, through
+#: ``_dependency_predicates`` -- which only knows how to read an inbound
+#: foreign key.
+_EXECUTABLE_CLASSIFICATIONS: frozenset[DependencyClassification] = frozenset(
+    {
+        DependencyClassification.OWNED,
+        DependencyClassification.ASSOCIATION,
+        DependencyClassification.BLOCK,
+    }
+)
+
+
 def _validate_executable_declarations(policy: PurgeEntityPolicy) -> None:
-    """Reject executable classifications missing their required action 
metadata."""
+    """Reject declarations the shared cleanup could not carry out.
+
+    Each shape here would otherwise raise mid-purge: for the scheduled task,
+    after the write-ahead audit row is already committed, so the attempt is
+    recorded as a failure on every run. Refused at registration instead.
+    """
+    metadata: sa.MetaData = sa.inspect(policy.model).local_table.metadata
     for dependency in policy.dependencies:
+        key: DependencyKey = dependency.key
+        if dependency.classification in _EXECUTABLE_CLASSIFICATIONS:
+            if key.kind != "foreign_key" or key.direction != "inbound":
+                raise RuntimeError(
+                    f"{dependency.classification.value} dependency "
+                    f"{key.describe()} is not an inbound foreign key; the "
+                    "shared cleanup can only act on those"
+                )
+            if key.related_table not in metadata.tables:
+                raise RuntimeError(
+                    f"{dependency.classification.value} dependency "
+                    f"{key.describe()} names a table outside the root's "
+                    "metadata, which cleanup cannot resolve"
+                )
         if (
-            dependency.classification is 
DependencyClassification.LISTENER_EFFECT
-            and dependency.listener_action is None
+            dependency.classification is DependencyClassification.BLOCK
+            and dependency.blocker is None
         ):
-            raise RuntimeError(
-                f"Missing listener action for {dependency.key.describe()}"
-            )
+            raise RuntimeError(f"Missing blocker reason for {key.describe()}")
+        if dependency.classification is 
DependencyClassification.LISTENER_EFFECT:
+            if dependency.listener_action is None:
+                raise RuntimeError(
+                    f"Missing listener action for {dependency.key.describe()}"
+                )
+            if dependency.phase is None:

Review Comment:
   A graph-complete listener declaration with `DELETE_DATASET_PERMISSION` at 
`ExecutionPhase.OWNED` passes admission, but stock callbacks dispatch listener 
effects only in `ASSOCIATIONS` and `POST_DELETE`, so the root can be purged 
without its declared cleanup ever running. Could admission validate the 
action/phase combinations the selected callbacks actually execute?



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