sadpandajoe commented on code in PR #44892:
URL: https://github.com/apache/superset/pull/44892#discussion_r4186369948
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -905,32 +1080,311 @@ 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 _validate_executable_declarations(policy: PurgeEntityPolicy) -> None:
- """Reject executable classifications missing their required action
metadata."""
+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
DependencyClassification.LISTENER_EFFECT
- and dependency.listener_action is None
- ):
+ 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"Missing listener action for {dependency.key.describe()}"
+ 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"
)
- if (
- dependency.classification is DependencyClassification.VERSION_OWNED
- and dependency.version_column is None
- ):
+
+
+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"
+ )
+
+
+#: The phase each listener action is dispatched in by the stock callbacks:
+#: tag cleanup runs with the association deletes, permission cleanup once the
+#: entity row is gone and the captured permission name is available.
+_LISTENER_ACTION_PHASES: Mapping[ListenerAction, ExecutionPhase] =
MappingProxyType(
+ {
+ ListenerAction.DELETE_TAGGED_OBJECTS: ExecutionPhase.ASSOCIATIONS,
+ ListenerAction.DELETE_DATASET_PERMISSION: ExecutionPhase.POST_DELETE,
+ }
+)
+
+
+#: 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_listener_dispatch(
+ policy: PurgeEntityPolicy, dependency: DependencyPolicy
+) -> None:
+ """Reject a listener effect no selected callback would dispatch.
+
+ Effects are dispatched from inside two stock callbacks -- association
+ cleanup and permission cleanup -- each for one phase. A declaration in any
+ other phase, or in a phase whose callback this policy replaces, is never
+ executed: the root purges and its declared cleanup silently does not run.
+ """
+ action: ListenerAction = cast(ListenerAction, dependency.listener_action)
+ expected_phase: ExecutionPhase | None = _LISTENER_ACTION_PHASES.get(action)
+ if expected_phase is None:
+ raise RuntimeError(
+ f"Unsupported listener action for {dependency.key.describe()}"
+ )
+ if dependency.phase is not expected_phase:
+ raise RuntimeError(
+ f"Listener effect {dependency.key.describe()} is declared for "
+ f"{cast(ExecutionPhase, dependency.phase).value}, but "
+ f"{action.value} is dispatched only in {expected_phase.value}"
+ )
+ if expected_phase is ExecutionPhase.ASSOCIATIONS:
+ dispatched = policy.delete_associations is delete_associations
+ callback = "delete_associations"
+ else:
+ dispatched = policy.cleanup_permission is cleanup_dataset_permission
+ callback = "cleanup_permission"
+ if not dispatched:
+ raise RuntimeError(
+ f"Listener effect {dependency.key.describe()} is dispatched by the
"
+ f"stock {callback} callback, which this policy replaces"
+ )
+
+
+def _validate_dependency_declaration(
+ policy: PurgeEntityPolicy,
+ dependency: DependencyPolicy,
+ metadata: sa.MetaData,
+) -> None:
+ """Reject one declaration the shared cleanup could not carry out."""
+ 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"Missing version target column for
{dependency.key.describe()}"
+ 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.BLOCK
+ and dependency.blocker is None
+ ):
+ 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 {key.describe()}")
+ if dependency.phase is None:
+ # Listener effects are dispatched by phase, so one without a phase
+ # matches none and would never run.
+ raise RuntimeError(f"Missing execution phase for {key.describe()}")
+ _validate_listener_dispatch(policy, dependency)
+ if (
+ dependency.classification is DependencyClassification.VERSION_OWNED
+ and dependency.version_column is None
+ ):
Review Comment:
A versioned `workflow → stage → leaf` policy can declare
`leaf_version.stage_id`, but version cleanup compares that column directly with
the workflow ID: purging workflow 1 with stage 10 deletes history for a
surviving workflow’s stage 1 and retains stage 10’s history. Could admission
reject intermediate-owner version targets until cleanup builds root-relative
predicates?
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -860,16 +894,155 @@ 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)
-@lru_cache(maxsize=None)
-def _validated_purge_policy(model: type[Any]) -> PurgeEntityPolicy:
- """Validate and return one root policy without blocking unrelated roots."""
- try:
- policy: PurgeEntityPolicy = purge_policy_registry()[model]
- except KeyError as ex:
- raise ValueError(f"Unsupported purge model: {model.__name__}") from ex
+#: Sentinel distinguishing "no index resolved yet" from an index resolved for
+#: no provider, which is the ordinary case.
+_UNRESOLVED: Any = object()
+
+
+@dataclass
+class _ResolvedRegistry:
+ """The index built for one provider, with the roots validated so far.
+
+ Compared by provider *identity* rather than memoized with ``lru_cache``: a
+ cache keyed on the provider would hash it, and a host callable
+ implemented as a mutable dataclass instance is unhashable. That would
+ raise here -- before the host boundary could isolate the host's mistake --
+ and stop the built-in roots purging too.
+ """
+
+ provider: Any = _UNRESOLVED
+ registry: Mapping[type[Any], PurgeEntityPolicy] = field(
+ default_factory=lambda: MappingProxyType({})
+ )
+ validated: set[type[Any]] = field(default_factory=set)
+
+
+_RESOLVED: _ResolvedRegistry = _ResolvedRegistry()
+
+
+def purge_policy_registry() -> Mapping[type[Any], PurgeEntityPolicy]:
+ """Index the built-in purge roots plus any the host installed.
+
+ Rebuilt only when the installed provider changes. Freezing 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, while resolving on every call would invoke the provider once per
+ root and could hand back different objects each time.
+ """
+ provider: Callable[[], Any] | None = (
+ current_app.config.get(HOST_POLICIES_CONFIG_KEY) if has_app_context()
else None
+ )
+ if _RESOLVED.provider is not provider:
+ host_policies: tuple[PurgeEntityPolicy, ...] =
_host_purge_policies(provider)
Review Comment:
Two concurrent first-use requests can invoke the same provider: after one
caches valid host policies, a slower invocation that fails transiently
overwrites them with the built-in-only fallback and caches that provider
identity, leaving host roots unsupported until restart or replacement. Could
registry initialization be single-flight and publish its resolved state
atomically?
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -860,16 +894,155 @@ 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__,
Review Comment:
If a host mistakenly supplies a SQLAlchemy `Table` as a policy’s `model`,
validation rejects it but this logging access raises `AttributeError` outside
the guard; registry initialization then fails for every built-in root too.
Could malformed model values be validated or safely named inside the
host-isolation boundary?
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -905,32 +1080,311 @@ 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 _validate_executable_declarations(policy: PurgeEntityPolicy) -> None:
- """Reject executable classifications missing their required action
metadata."""
+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
DependencyClassification.LISTENER_EFFECT
- and dependency.listener_action is None
- ):
+ 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:
Review Comment:
A host root with primary key `(tenant_id, id)` and `deleted_at` passes these
checks, but the scheduled task calls `session.get(model, entity_id)` with a
scalar, so every eligible row fails while dry-run reports it as purgeable.
Could admission validate the mapper identity the stock task requires, or carry
the complete key through scanning and lookup?
##########
tests/unit_tests/tasks/test_deletion_retention.py:
##########
@@ -330,6 +331,79 @@ class UnsupportedModel(SoftDeleteMixin):
engine.dispose()
+def test_scan_failure_keeps_the_counts_earned_before_it(app_context: None) ->
None:
+ """A scan that fails part-way reports itself and keeps what it purged.
+
+ By the time a later page fails, the earlier page's deletions are
+ committed. Discarding the counts would make the run's summary and its
+ purge gauges understate what was actually removed.
+ """
+ import superset.tasks.deletion_retention as mod
+ from superset.commands.deletion_retention.purge_cascade import
CascadeResult
+ from superset.models.slice import Slice
+
+ purged_result: CascadeResult = CascadeResult(
+ purged=True, entity_type="chart", entity_uuid="gone"
+ )
+
+ def pages(*args: Any, **kwargs: Any) -> Iterator[list[int]]:
+ yield [1, 2]
+ raise RuntimeError("no such column: id")
+
+ with (
+ patch.object(mod, "_iter_eligible_ids", side_effect=pages),
+ patch.object(mod, "_purge_one", return_value=purged_result),
+ ):
+ purged, would, failures, blocked, scan_failures = mod._purge_model(
+ Slice, datetime.now(), dry_run=False
+ )
+
+ assert (purged, would, failures, blocked) == (2, 0, 0, 0)
+ assert scan_failures == 1
+
+
+def test_root_without_a_table_name_does_not_abort_the_run(
+ app_config: Config,
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ """Resolving a root's table name is itself guarded.
+
+ The name is read off the model, so a root that cannot supply one must be
+ counted and skipped like any other failing root -- not end the pass before
+ the roots that could have purged are reached.
+ """
+ # avoid app-init regression: model helpers require the app_config fixture
first.
+ from superset.models.helpers import SoftDeleteMixin
+ from superset.tasks import deletion_retention as task
+
+ supported_models: list[type[SoftDeleteMixin]] =
list(task.purge_policy_registry())
+
+ class NoTableName(SoftDeleteMixin):
Review Comment:
Defining this subclass mutates the original process-global registry before
`monkeypatch` saves it, so teardown restores a list still containing
`NoTableName`; later retention passes unexpectedly report a scan failure for
the tableless class. Could the registry be patched to a copied list before
defining the subclass, as the neighboring unsupported-model test does?
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -905,32 +1080,311 @@ 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 _validate_executable_declarations(policy: PurgeEntityPolicy) -> None:
- """Reject executable classifications missing their required action
metadata."""
+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
DependencyClassification.LISTENER_EFFECT
- and dependency.listener_action is None
- ):
+ 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"Missing listener action for {dependency.key.describe()}"
+ 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"
)
- if (
- dependency.classification is DependencyClassification.VERSION_OWNED
- and dependency.version_column is None
- ):
+
+
+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"
+ )
+
+
+#: The phase each listener action is dispatched in by the stock callbacks:
+#: tag cleanup runs with the association deletes, permission cleanup once the
+#: entity row is gone and the captured permission name is available.
+_LISTENER_ACTION_PHASES: Mapping[ListenerAction, ExecutionPhase] =
MappingProxyType(
+ {
+ ListenerAction.DELETE_TAGGED_OBJECTS: ExecutionPhase.ASSOCIATIONS,
+ ListenerAction.DELETE_DATASET_PERMISSION: ExecutionPhase.POST_DELETE,
+ }
+)
+
+
+#: 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_listener_dispatch(
+ policy: PurgeEntityPolicy, dependency: DependencyPolicy
+) -> None:
+ """Reject a listener effect no selected callback would dispatch.
+
+ Effects are dispatched from inside two stock callbacks -- association
+ cleanup and permission cleanup -- each for one phase. A declaration in any
+ other phase, or in a phase whose callback this policy replaces, is never
+ executed: the root purges and its declared cleanup silently does not run.
+ """
+ action: ListenerAction = cast(ListenerAction, dependency.listener_action)
+ expected_phase: ExecutionPhase | None = _LISTENER_ACTION_PHASES.get(action)
+ if expected_phase is None:
+ raise RuntimeError(
+ f"Unsupported listener action for {dependency.key.describe()}"
+ )
+ if dependency.phase is not expected_phase:
+ raise RuntimeError(
+ f"Listener effect {dependency.key.describe()} is declared for "
+ f"{cast(ExecutionPhase, dependency.phase).value}, but "
+ f"{action.value} is dispatched only in {expected_phase.value}"
+ )
+ if expected_phase is ExecutionPhase.ASSOCIATIONS:
+ dispatched = policy.delete_associations is delete_associations
+ callback = "delete_associations"
+ else:
+ dispatched = policy.cleanup_permission is cleanup_dataset_permission
Review Comment:
A host permission listener declared at `POST_DELETE` with this stock
callback passes admission, yet retaining `dataset_permission_name` for a
non-dataset entity always captures `None`; the cascade skips
`cleanup_permission` entirely and deletes the root without its declared
permission cleanup. Could admission also check the capture prerequisite, or
fail closed when a required effect has no captured name?
##########
superset/commands/deletion_retention/purge_policy.py:
##########
@@ -905,32 +1080,311 @@ 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 _validate_executable_declarations(policy: PurgeEntityPolicy) -> None:
- """Reject executable classifications missing their required action
metadata."""
+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
DependencyClassification.LISTENER_EFFECT
- and dependency.listener_action is None
- ):
+ 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"Missing listener action for {dependency.key.describe()}"
+ 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"
)
- if (
- dependency.classification is DependencyClassification.VERSION_OWNED
- and dependency.version_column is None
- ):
+
+
+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"
+ )
+
+
+#: The phase each listener action is dispatched in by the stock callbacks:
+#: tag cleanup runs with the association deletes, permission cleanup once the
+#: entity row is gone and the captured permission name is available.
+_LISTENER_ACTION_PHASES: Mapping[ListenerAction, ExecutionPhase] =
MappingProxyType(
+ {
+ ListenerAction.DELETE_TAGGED_OBJECTS: ExecutionPhase.ASSOCIATIONS,
+ ListenerAction.DELETE_DATASET_PERMISSION: ExecutionPhase.POST_DELETE,
+ }
+)
+
+
+#: 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_listener_dispatch(
+ policy: PurgeEntityPolicy, dependency: DependencyPolicy
+) -> None:
+ """Reject a listener effect no selected callback would dispatch.
+
+ Effects are dispatched from inside two stock callbacks -- association
+ cleanup and permission cleanup -- each for one phase. A declaration in any
+ other phase, or in a phase whose callback this policy replaces, is never
+ executed: the root purges and its declared cleanup silently does not run.
+ """
+ action: ListenerAction = cast(ListenerAction, dependency.listener_action)
+ expected_phase: ExecutionPhase | None = _LISTENER_ACTION_PHASES.get(action)
+ if expected_phase is None:
+ raise RuntimeError(
+ f"Unsupported listener action for {dependency.key.describe()}"
+ )
+ if dependency.phase is not expected_phase:
+ raise RuntimeError(
+ f"Listener effect {dependency.key.describe()} is declared for "
+ f"{cast(ExecutionPhase, dependency.phase).value}, but "
+ f"{action.value} is dispatched only in {expected_phase.value}"
+ )
+ if expected_phase is ExecutionPhase.ASSOCIATIONS:
+ dispatched = policy.delete_associations is delete_associations
+ callback = "delete_associations"
+ else:
+ dispatched = policy.cleanup_permission is cleanup_dataset_permission
+ callback = "cleanup_permission"
+ if not dispatched:
+ raise RuntimeError(
+ f"Listener effect {dependency.key.describe()} is dispatched by the
"
+ f"stock {callback} callback, which this policy replaces"
+ )
+
+
+def _validate_dependency_declaration(
+ policy: PurgeEntityPolicy,
+ dependency: DependencyPolicy,
+ metadata: sa.MetaData,
+) -> None:
+ """Reject one declaration the shared cleanup could not carry out."""
+ 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"Missing version target column for
{dependency.key.describe()}"
+ f"{dependency.classification.value} dependency "
+ f"{key.describe()} names a table outside the root's "
Review Comment:
With both default-schema `root`/`child` and `host.root`/`host.child` in the
metadata, a policy for `host.root` passes admission but stock cleanup resolves
the unqualified names to the default-schema tables; purging host root 1 can
delete children of unrelated live default-schema roots. Could schema-qualified
table identities be preserved throughout discovery and execution, or this shape
be rejected before admission?
--
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]