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


##########
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:
   <!-- Bito Reply -->
   The suggestion to replace `datetime.utcnow()` with 
`datetime.now(timezone.utc)` is generally recommended to ensure timezone-aware 
datetimes. However, since the project intentionally uses naive datetimes to 
match the metadata database timestamp representation, the current 
implementation is appropriate for this specific context. You may choose to 
ignore this suggestion if maintaining consistency with the database schema is 
the priority.



##########
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:
   <!-- Bito Reply -->
   The suggestion to add `tzinfo=timezone.utc` is not appropriate in this 
specific context. As you noted, `normalize_window` is designed to compare aware 
timestamps against plain database `DateTime` columns by stripping timezone 
information and microseconds. Adding `tzinfo` would create a mismatch with the 
function's intended contract for these comparisons.
   
   **tests/unit_tests/commands/report/execution_claim_test.py**
   ```
   ) == datetime(2026, 9, 15, tzinfo=timezone.utc)
   ```



##########
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:
   <!-- Bito Reply -->
   The suggestion to use named constants instead of hardcoded strings for 
exception names is appropriate. It improves maintainability by ensuring that if 
an exception class is renamed, the check remains accurate without requiring 
manual updates to string literals. You should apply this change by defining 
constants for these exception names and referencing them in the check.
   
   **superset/commands/report/execute.py**
   ```
   if terminal_reason not in (
                   SoftTimeLimitExceeded.__name__,
                   ReportExecutionBudgetExceededError.__name__,
               ):
   ```



##########
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:
   <!-- Bito Reply -->
   The suggestion to replace `datetime.utcnow()` with 
`datetime.now(timezone.utc).replace(tzinfo=None)` is appropriate. Using 
`datetime.now(timezone.utc)` is the recommended practice for generating 
timezone-aware timestamps, and replacing the timezone info ensures consistency 
with the existing codebase in `superset/commands/report/execute.py` while 
avoiding potential precision issues associated with `utcnow()`.
   
   **superset/commands/report/execute.py**
   ```
   ReportSchedule.last_state: ReportState.ERROR,
                           ReportSchedule.last_eval_dttm: 
datetime.now(timezone.utc).replace(tzinfo=None),
   ```



##########
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:
   <!-- Bito Reply -->
   The suggestion provided by the reviewer is appropriate and addresses a 
potential TypeError. By ensuring that `entered_at` is only calculated when 
`last_eval_dttm` is available or when `last_working` is present, the code 
avoids performing arithmetic on `None` values. Applying this change improves 
the robustness of the timeout check.
   
   **superset/commands/report/execute.py**
   ```
   entered_at = self._report_schedule.last_eval_dttm
           if not self._report_schedule.execution_owner and last_working:
               entered_at = last_working.end_dttm
           if entered_at is None:
               return False
   ```



##########
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:
   <!-- Bito Reply -->
   The suggestion from the reviewer is to ensure that when `_delivery_started` 
is True, the schedule state is terminalized before returning `False` in 
`_handle_retry_or_error`. Currently, the code resets the retry counter and 
returns `False`, which may leave the schedule in a `WORKING` state with an 
`execution_owner` set if the caller's subsequent error handling also fails or 
is skipped. Persisting a terminal state would prevent the schedule from 
remaining stuck until the working-timeout reaper intervenes.
   
   **superset/commands/report/execute.py**
   ```
   +        if self._delivery_started:
   +            self._reset_retry_counter()
   +            return False
   ```



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