msyavuz commented on code in PR #44892:
URL: https://github.com/apache/superset/pull/44892#discussion_r4183383151
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -860,7 +894,96 @@ def declare(
cleanup_permission=cleanup_dataset_permission,
),
}
- return validate_unique_root_policies(registry.values())
+ return tuple(registry.values())
+
+
+def _host_purge_policies() -> 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: an unavailable provider, a malformed payload, a policy that
+ collides with a root declared here, and two policies for one host root are
+ each logged and dropped. A broken host declaration must not take the
+ scheduled purge down with it, and must never redefine how a chart,
+ dashboard or dataset is purged.
+ """
+ if not has_app_context():
+ return ()
+ provider: Callable[[], Any] | None = current_app.config.get(
+ HOST_POLICIES_CONFIG_KEY
+ )
+ 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
+ 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.
+
+ Uncached on purpose, unlike its two inputs: the built-in declarations are
+ built once per process, while a host policy is resolved per call, so a
+ provider installed after the first purge is honored for any root not yet
+ resolved. ``get_purge_policy`` memoizes per model, so replacing the policy
+ of a root it has already resolved needs a restart -- and that memoization
+ is why rebuilding this small index is off the hot path.
+ """
+ return validate_unique_root_policies(
Review Comment:
**major**: purge_policy_registry() is no longer cached and now re-runs
provider() and validate_unique_root_policies on every call, and `_purge_roots`
calls it once per model. A slow or non-deterministic provider is invoked
repeatedly, and the returned policy objects can differ between calls.
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -860,7 +894,96 @@ def declare(
cleanup_permission=cleanup_dataset_permission,
),
}
- return validate_unique_root_policies(registry.values())
+ return tuple(registry.values())
+
+
+def _host_purge_policies() -> 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: an unavailable provider, a malformed payload, a policy that
+ collides with a root declared here, and two policies for one host root are
+ each logged and dropped. A broken host declaration must not take the
+ scheduled purge down with it, and must never redefine how a chart,
+ dashboard or dataset is purged.
+ """
+ if not has_app_context():
Review Comment:
**major**: Outside an app context, host roots silently drop out and then get
reported as unsupported. Because `get_purge_policy` memoizes, a failure from a
context-less first call may be cached and persist until restart.
##########
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
Review Comment:
**major**: The broad except now also swallows a RuntimeError from policy
validation (`Incomplete purge policy`) raised via `purge_policy_registry` or
`get_purge_policy`. A bad host policy then shows up only as `scan_failures` on
every run, and it is mislabeled as a scan error.
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -909,23 +1032,211 @@ def _validated_purge_policy(model: type[Any]) ->
PurgeEntityPolicy:
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:
Review Comment:
**minor**: `_tables_share_a_foreign_key` only compares the first element of
each foreign key constraint and only checks direct FKs. Cross-references
between owned siblings through a third table or a multi-column FK are not
detected, so the sibling-order check can miss real cases.
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -329,34 +363,34 @@ def _discover_table_dependencies(
if constraint.elements[0].column.table is table:
discovered.add(_fk_key(table, constraint, "inbound"))
discovered.update(
- _discover_recursive_table_dependencies(
- table,
- recursive_tables=recursive_tables,
- visited=next_visited,
- )
+ _discover_recursive_table_dependencies(table, recursive_tables, seen)
)
return frozenset(discovered)
def _discover_recursive_table_dependencies(
table: sa.Table,
- *,
recursive_tables: frozenset[str],
- visited: frozenset[str],
+ seen: _Traversal,
) -> frozenset[DependencyKey]:
- """Discover dependencies for named recursive tables not already visited."""
+ """Discover dependencies for named recursive tables not already visited.
+
+ Sorted so a traversal is reproducible, and driven off the shared
+ bookkeeping so each table is walked at most once per kind of walk.
+ """
discovered: set[DependencyKey] = set()
- for recursive_table_name in recursive_tables - visited:
+ for recursive_table_name in sorted(recursive_tables - seen.tables):
+ # Re-checked inside the loop: the set grows as the walk descends, so
+ # a table a nested call already covered is skipped rather than
+ # re-entered to return nothing.
+ if recursive_table_name in seen.tables:
Review Comment:
**minor**: The shared `seen` set changes discovery output for graphs where a
table is reached via different paths. The `_Traversal` split of tables and
mappers was added to mitigate this, but I only see tests for the call count,
not for built-in key-set equivalence.
--
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]