fitzee commented on code in PR #44336:
URL: https://github.com/apache/superset/pull/44336#discussion_r4060270402
##########
UPDATING.md:
##########
@@ -24,11 +24,20 @@ assists people when migrating to a new version.
## Next
-- With `SEMANTIC_LAYERS` enabled, combined connection discovery honors
`Database.can_read` and `SemanticLayer.can_read` independently. Each permitted
source retains its normal row filters, including dynamic database filters for
Admin. A source filter never includes rows or counts from a denied source;
callers with neither read permission are denied. Feature-off database browsing
is unchanged.
-- The combined datasource list (`GET /api/v1/datasource/`) accepts Dataset
read without an additional Datasource read grant, regardless of
`SEMANTIC_LAYERS`. With the flag enabled, SemanticView read independently
permits semantic-view discovery. Existing row-level dataset/chart access
remains enforced.
-- The `presto` extra requires PyHive 0.7.0 or later. PyHive 0.6.5 cannot load
- its Presto dialect under SQLAlchemy 2 because it imports
`sqlalchemy.databases`.
- Upgrade existing installations with `pip install "pyhive[presto]>=0.7.0"`.
+### Scheduled report and alert retry admission
+
+Run `superset db upgrade` before starting workers with this version. The
migration
+adds nullable `execution_owner` and `execution_window` columns to
`report_schedule`.
+Pause scheduling and drain in-flight executions and queued retry tasks before
+upgrading workers together: older workers do not participate in execution
fencing.
+Retries queued by the old task signature are discarded rather than replayed
without
+ownership evidence. Restart scheduling after migration and worker replacement.
Review Comment:
Corrected in ba4ae23ba3a37f74ad7b09ecba5b46c182982565. The previous wording
was too broad: two-argument legacy retries are identified and discarded, but
older one-argument retries are indistinguishable from fresh cron tasks.
UPDATING.md explicitly requires draining those messages and no longer promises
they will be discarded automatically.
##########
tests/unit_tests/utils/test_screenshot_utils.py:
##########
@@ -430,6 +475,147 @@ def
test_report_mode_rejects_partial_first_tile_fallback(self):
class TestTakeTiledScreenshot:
+ def test_per_tile_diagnostics_failure_does_not_discard_capture(self,
mock_page):
+ """Diagnostics must not discard valid tiles: a non-timeout evaluate
+ failure on the per-tile diagnostics path logs a warning and the
+ capture still succeeds."""
+ final_states = [{"chartId": "1", "state": "rendered"}]
+ diagnostics_calls = 0
+
+ def evaluate(script, _arg=None):
+ nonlocal diagnostics_calls
+ if "scrollWidth" in script:
+ return {"height": 1000, "top": 100, "left": 50, "width": 800}
+ if script == CONTENTFUL_CHART_HOLDERS_IN_CLIP_JS:
+ return {"total": 1, "contentful": 1}
+ if script == FIND_CHART_HOLDER_STATES_JS:
+ diagnostics_calls += 1
+ if diagnostics_calls == 1:
+ raise RuntimeError("Execution context was destroyed")
+ return final_states
+ return None
+
+ mock_page.wait_for_function.side_effect = None
+ mock_page.wait_for_function.return_value = None
+ mock_page.evaluate.side_effect = evaluate
+ combined = self._create_chart_like_tile()
+ mock_page.screenshot.return_value = combined
+
+ with patch("superset.utils.screenshot_utils.logger") as mock_logger:
+ with patch(
+ "superset.utils.screenshot_utils.combine_screenshot_tiles",
+ return_value=combined,
+ ):
+ result = take_tiled_screenshot(
+ mock_page,
+ "dashboard",
+ tile_height=2000,
+ load_wait=30,
+ report_execution_context=_report_context(),
+ )
+
+ assert result == combined
+ assert mock_page.screenshot.call_count == 1
+ assert any(
+ call.args
+ and call.args[0].startswith(
+ "Unable to collect per-tile chart-holder diagnostics"
+ )
+ for call in mock_logger.warning.call_args_list
+ )
+
+ def test_final_semantic_status_fires_exactly_once_for_error_holders(
+ self, mock_page
+ ):
+ """One error chart spanning every tile emits ONE report_semantic_status
+ WARNING for the capture (final block), not one per tile — the per-tile
+ INFO lines already carry error_holders."""
+ with_error = [
+ {"chartId": "1", "state": "rendered"},
+ {"chartId": "2", "state": "error"},
+ ]
+ mock_page.wait_for_function.side_effect = None
+ mock_page.wait_for_function.return_value = None
+ mock_page.evaluate.side_effect = [
+ {"height": 5000, "top": 100, "left": 50, "width": 800},
+ None,
+ with_error,
+ None,
+ with_error,
+ None,
+ with_error,
+ with_error,
Review Comment:
This capture-validation test was removed from #44336 when that work was
split into #44465. The screenshot utilities and their tests are no longer
changed by this retry PR. No additional change is needed here.
##########
tests/unit_tests/commands/report/execution_claim_test.py:
##########
@@ -0,0 +1,322 @@
+# 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.
+
+import os
+from collections.abc import Iterator
+from concurrent.futures import ThreadPoolExecutor
+from datetime import datetime, timedelta, timezone
+from pathlib import Path
+from threading import Barrier
+from typing import Any
+from unittest.mock import Mock
+from uuid import uuid4
+
+import pytest
+import sqlalchemy as sa
+from pytest_mock import MockerFixture
+from sqlalchemy.orm import sessionmaker
+
+from superset.commands.report.execution_claim import claim_execution,
normalize_window
+from superset.reports.models import ReportSchedule, ReportScheduleType,
ReportState
+
+
[email protected](params=["sqlite", "postgres"])
+def claim_engine(tmp_path: Path, request: pytest.FixtureRequest) ->
Iterator[sa.Engine]:
+ """Exercise PostgreSQL REPEATABLE READ when a test database is provided."""
+ if request.param == "sqlite":
+ engine = sa.create_engine(f"sqlite:///{tmp_path / 'claims.db'}")
+ ReportSchedule.__table__.create(engine)
+ yield engine
+ engine.dispose()
+ return
+ uri = os.environ.get("REPORT_CLAIM_TEST_DATABASE_URI") or os.environ.get(
+ "SUPERSET__SQLALCHEMY_DATABASE_URI"
+ )
+ if not uri or sa.engine.make_url(uri).get_backend_name() != "postgresql":
+ pytest.skip("Set REPORT_CLAIM_TEST_DATABASE_URI for PostgreSQL race
coverage")
+ schema = f"report_claim_test_{uuid4().hex}"
+ base_engine = sa.create_engine(uri, isolation_level="REPEATABLE READ")
+ table = ReportSchedule.__table__.to_metadata(sa.MetaData(), schema=schema)
+ with base_engine.begin() as connection:
+ connection.execute(sa.schema.CreateSchema(schema))
+ connection.execute(
+ sa.schema.CreateTable(table, include_foreign_key_constraints=[])
+ )
+ engine = base_engine.execution_options(schema_translate_map={None: schema})
+ try:
+ yield engine
+ finally:
+ with base_engine.begin() as connection:
+ connection.execute(sa.schema.DropSchema(schema, cascade=True))
+ base_engine.dispose()
+
+
[email protected](params=[ReportScheduleType.REPORT, ReportScheduleType.ALERT])
+def sessions(claim_engine: sa.Engine, request: pytest.FixtureRequest) ->
sessionmaker:
+ """Use real independent connections, not a shared mocked session."""
+ factory = sessionmaker(bind=claim_engine)
+ with factory() as session:
+ session.add(
+ ReportSchedule(
+ id=1,
+ name="claim test",
+ type=request.param,
+ crontab="* * * * *",
+ last_state=ReportState.NOOP,
+ retry_on_failure=True,
+ )
+ )
+ session.commit()
+ return factory
+
+
+def claim(session, window, *, retry=False, owner=None):
+ return claim_execution(
+ session,
+ 1,
+ str(uuid4()),
+ window,
+ is_retry=retry,
+ expected_owner=owner,
+ retries_enabled=True,
+ stale_retry_seconds=4140,
+ )
+
+
[email protected]("retry", [False, True])
+def test_competing_workers_have_one_winner(sessions, retry):
+ window = datetime.utcnow().replace(microsecond=0)
+ owner = str(uuid4())
+ if retry:
+ with sessions() as session:
+ session.query(ReportSchedule).update(
+ {
+ "last_state": ReportState.RETRYING,
+ "execution_owner": owner,
+ "execution_window": window,
+ "retry_scheduled_dttm": window,
+ }
+ )
+ session.commit()
+ # Synchronize AFTER both SELECTs. Each worker must contend on the UPDATE
+ # using the same observed state, rather than simply reading the winner.
+ barrier = Barrier(2)
+ engine = sessions.kw["bind"]
+
+ def synchronize(_conn, _cursor, statement, _parameters, _context,
_executemany):
+ if statement.startswith("UPDATE ") and "report_schedule" in statement:
+ barrier.wait(timeout=10)
+
+ sa.event.listen(engine, "before_cursor_execute", synchronize)
+ try:
+
+ def run(_index):
+ with sessions() as session:
+ return claim(session, window, retry=retry, owner=owner)
+
+ with ThreadPoolExecutor(max_workers=2) as workers:
+ results = list(workers.map(run, range(2)))
+ assert sum(result is not None for result in results) == 1
+ finally:
+ sa.event.remove(engine, "before_cursor_execute", synchronize)
+
+
+def test_retry_window_and_owner_fence_replays(sessions):
+ window = datetime.utcnow().replace(microsecond=0)
Review Comment:
Updated in ba4ae23ba3a37f74ad7b09ecba5b46c182982565. The claim tests use
`datetime.now(timezone.utc).replace(tzinfo=None)` through a helper. Keeping the
final value naive is intentional: it matches the metadata database timestamp
representation.
##########
tests/unit_tests/commands/report/execution_claim_test.py:
##########
@@ -0,0 +1,322 @@
+# 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.
+
+import os
+from collections.abc import Iterator
+from concurrent.futures import ThreadPoolExecutor
+from datetime import datetime, timedelta, timezone
+from pathlib import Path
+from threading import Barrier
+from typing import Any
+from unittest.mock import Mock
+from uuid import uuid4
+
+import pytest
+import sqlalchemy as sa
+from pytest_mock import MockerFixture
+from sqlalchemy.orm import sessionmaker
+
+from superset.commands.report.execution_claim import claim_execution,
normalize_window
+from superset.reports.models import ReportSchedule, ReportScheduleType,
ReportState
+
+
[email protected](params=["sqlite", "postgres"])
+def claim_engine(tmp_path: Path, request: pytest.FixtureRequest) ->
Iterator[sa.Engine]:
+ """Exercise PostgreSQL REPEATABLE READ when a test database is provided."""
+ if request.param == "sqlite":
+ engine = sa.create_engine(f"sqlite:///{tmp_path / 'claims.db'}")
+ ReportSchedule.__table__.create(engine)
+ yield engine
+ engine.dispose()
+ return
+ uri = os.environ.get("REPORT_CLAIM_TEST_DATABASE_URI") or os.environ.get(
+ "SUPERSET__SQLALCHEMY_DATABASE_URI"
+ )
+ if not uri or sa.engine.make_url(uri).get_backend_name() != "postgresql":
+ pytest.skip("Set REPORT_CLAIM_TEST_DATABASE_URI for PostgreSQL race
coverage")
+ schema = f"report_claim_test_{uuid4().hex}"
+ base_engine = sa.create_engine(uri, isolation_level="REPEATABLE READ")
+ table = ReportSchedule.__table__.to_metadata(sa.MetaData(), schema=schema)
+ with base_engine.begin() as connection:
+ connection.execute(sa.schema.CreateSchema(schema))
+ connection.execute(
+ sa.schema.CreateTable(table, include_foreign_key_constraints=[])
+ )
+ engine = base_engine.execution_options(schema_translate_map={None: schema})
+ try:
+ yield engine
+ finally:
+ with base_engine.begin() as connection:
+ connection.execute(sa.schema.DropSchema(schema, cascade=True))
+ base_engine.dispose()
+
+
[email protected](params=[ReportScheduleType.REPORT, ReportScheduleType.ALERT])
+def sessions(claim_engine: sa.Engine, request: pytest.FixtureRequest) ->
sessionmaker:
+ """Use real independent connections, not a shared mocked session."""
+ factory = sessionmaker(bind=claim_engine)
+ with factory() as session:
+ session.add(
+ ReportSchedule(
+ id=1,
+ name="claim test",
+ type=request.param,
+ crontab="* * * * *",
+ last_state=ReportState.NOOP,
+ retry_on_failure=True,
+ )
+ )
+ session.commit()
+ return factory
+
+
+def claim(session, window, *, retry=False, owner=None):
+ return claim_execution(
+ session,
+ 1,
+ str(uuid4()),
+ window,
+ is_retry=retry,
+ expected_owner=owner,
+ retries_enabled=True,
+ stale_retry_seconds=4140,
+ )
+
+
[email protected]("retry", [False, True])
+def test_competing_workers_have_one_winner(sessions, retry):
+ window = datetime.utcnow().replace(microsecond=0)
+ owner = str(uuid4())
+ if retry:
+ with sessions() as session:
+ session.query(ReportSchedule).update(
+ {
+ "last_state": ReportState.RETRYING,
+ "execution_owner": owner,
+ "execution_window": window,
+ "retry_scheduled_dttm": window,
+ }
+ )
+ session.commit()
+ # Synchronize AFTER both SELECTs. Each worker must contend on the UPDATE
+ # using the same observed state, rather than simply reading the winner.
+ barrier = Barrier(2)
+ engine = sessions.kw["bind"]
+
+ def synchronize(_conn, _cursor, statement, _parameters, _context,
_executemany):
+ if statement.startswith("UPDATE ") and "report_schedule" in statement:
+ barrier.wait(timeout=10)
+
+ sa.event.listen(engine, "before_cursor_execute", synchronize)
+ try:
+
+ def run(_index):
+ with sessions() as session:
+ return claim(session, window, retry=retry, owner=owner)
+
+ with ThreadPoolExecutor(max_workers=2) as workers:
+ results = list(workers.map(run, range(2)))
+ assert sum(result is not None for result in results) == 1
+ finally:
+ sa.event.remove(engine, "before_cursor_execute", synchronize)
+
+
+def test_retry_window_and_owner_fence_replays(sessions):
+ window = datetime.utcnow().replace(microsecond=0)
+ with sessions() as session:
+ assert claim(session, window) is not None
+ schedule = session.query(ReportSchedule).one()
+ owner = schedule.execution_owner
+ schedule.last_state = ReportState.RETRYING
+ schedule.retry_scheduled_dttm = window
+ session.commit()
+ with sessions() as session:
+ assert claim(session, window) is None
+ assert claim(session, window, retry=True, owner="wrong-owner") is None
+ assert claim(session, window, retry=True, owner=owner) is not None
+ schedule = session.query(ReportSchedule).one()
+ schedule.last_state = ReportState.RETRYING
+ session.commit()
+ assert claim(session, window, retry=True, owner=owner) is None
+ schedule.last_state = ReportState.SUCCESS
+ session.commit()
+ assert claim(session, window) is None
+ assert claim(session, window + timedelta(hours=1)) is not None
+
+
+def test_alert_can_claim_a_retry(sessions):
+ window = datetime.utcnow()
+ with sessions() as session:
+ session.query(ReportSchedule).update(
+ {
+ "type": ReportScheduleType.ALERT,
+ "last_state": ReportState.RETRYING,
+ "execution_owner": "owner",
+ "retry_scheduled_dttm": window,
+ }
+ )
+ session.commit()
+ assert claim(session, window, retry=True, owner="owner") is not None
+
+
+def test_window_normalization_converts_offsets():
+ assert normalize_window(
+ datetime(2026, 9, 15, 10, tzinfo=timezone(timedelta(hours=10)))
+ ) == datetime(2026, 9, 15)
Review Comment:
Leaving this expectation unchanged. `normalize_window` deliberately converts
aware timestamps to UTC and removes timezone information and microseconds for
comparison with plain database DateTime columns. An aware expected value would
test the opposite contract.
##########
superset/commands/report/execute.py:
##########
@@ -215,6 +204,66 @@ def persist_owned_report_execution_terminal_error(
# this boundary. Roll back again so a failed terminal flush cannot
leave
# the scoped session unusable for the retry.
db.session.rollback() # pylint: disable=consider-using-transaction
+ if report_context is not None and report_context.execution_claimed:
+ schedule = (
+ db.session.query(ReportSchedule)
+ .filter_by(id=report_schedule_id)
+ .one_or_none()
+ )
+ if schedule is None:
+ return False
+ matched = (
+ db.session.query(ReportSchedule)
+ .filter(
+ ReportSchedule.id == report_schedule_id,
+ ReportSchedule.execution_owner == str(execution_id),
+ ReportSchedule.last_state == ReportState.WORKING,
+ )
+ .update(
+ {
+ ReportSchedule.last_state: ReportState.ERROR,
+ ReportSchedule.last_eval_dttm: datetime.utcnow(),
Review Comment:
Updated in ba4ae23ba3a37f74ad7b09ecba5b46c182982565 to
`datetime.now(timezone.utc).replace(tzinfo=None)`, consistent with the
surrounding code. The CAS remains protected by execution owner and WORKING
state; this change does not rely on a claimed precision difference between the
datetime APIs.
##########
superset/commands/report/execute.py:
##########
@@ -215,6 +204,66 @@ def persist_owned_report_execution_terminal_error(
# this boundary. Roll back again so a failed terminal flush cannot
leave
# the scoped session unusable for the retry.
db.session.rollback() # pylint: disable=consider-using-transaction
+ if report_context is not None and report_context.execution_claimed:
+ schedule = (
+ db.session.query(ReportSchedule)
+ .filter_by(id=report_schedule_id)
+ .one_or_none()
+ )
+ if schedule is None:
+ return False
+ matched = (
+ db.session.query(ReportSchedule)
+ .filter(
+ ReportSchedule.id == report_schedule_id,
+ ReportSchedule.execution_owner == str(execution_id),
+ ReportSchedule.last_state == ReportState.WORKING,
+ )
+ .update(
+ {
+ ReportSchedule.last_state: ReportState.ERROR,
+ ReportSchedule.last_eval_dttm: datetime.utcnow(),
+ },
+ synchronize_session="fetch",
+ )
+ )
+ if matched != 1:
+ db.session.rollback() # pylint:
disable=consider-using-transaction
+ return False
+ # The claim is durable even if the WORKING log was rolled back.
+ # create_log promotes that log when present and creates it
otherwise.
+ state = BaseReportState(
+ schedule,
+ schedule.execution_window,
+ execution_id,
+ report_context,
+ )
+ state.create_log(error_message, log_state=ReportState.ERROR)
+ logger.error(
+ "report_execution_terminal %s state=Error terminal_reason=%s",
+ report_context.log_context,
+ terminal_reason,
+ )
+ # Persist first; notification failure must not undo
terminalization.
+ # Timeout cleanup must not start another network operation.
+ if terminal_reason not in (
+ "SoftTimeLimitExceeded",
+ "ReportExecutionBudgetExceededError",
+ ):
Review Comment:
Updated in ba4ae23ba3a37f74ad7b09ecba5b46c182982565. The comparison uses
`SoftTimeLimitExceeded.__name__` and
`ReportExecutionBudgetExceededError.__name__` rather than independent string
literals.
##########
superset/commands/report/execute.py:
##########
@@ -1674,6 +1809,10 @@ def _handle_retry_or_error(
if not retry_on_failure:
return False
+ if self._delivery_started:
+ self._reset_retry_counter()
+ return False
Review Comment:
Leaving terminalization with the caller and command-boundary fallback.
Returning False is the helper's contract: the state handler persists ERROR, and
an escaping failure invokes the ownership-fenced fallback for both alerts and
reports. Existing tests cover partial delivery without retry and failures
escaping after a claim. If persistence itself remains unavailable, timeout
recovery is still necessary; persisting inside this helper would not remove
that database-failure case.
##########
superset/commands/report/execute.py:
##########
@@ -1741,8 +1880,13 @@ def is_on_working_timeout(self) -> bool:
last_working = ReportScheduleDAO.find_last_entered_working_log(
self._report_schedule
)
- if not last_working:
+ if not last_working and not self._report_schedule.execution_owner:
return False
+ # A worker can die between committing its claim and creating the log.
+ # In that gap, a previous execution's old log must not expire this
claim.
+ entered_at = self._report_schedule.last_eval_dttm
+ if not self._report_schedule.execution_owner and last_working:
+ entered_at = last_working.end_dttm
Review Comment:
Fixed in ba4ae23ba3a37f74ad7b09ecba5b46c182982565. The timeout predicate
checks the selected `entered_at` value, not only `last_eval_dttm`, before
comparing timestamps. Added a regression test for a legacy WORKING log with no
end timestamp.
--
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]