This is an automated email from the ASF dual-hosted git repository.
aminghadersohi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/superset.git
The following commit(s) were added to refs/heads/master by this push:
new 7029c468c4d fix(reports): bound CSV transport retries and sanitize
failures (#43977)
7029c468c4d is described below
commit 7029c468c4d58b40c695b49c4f2a7806a05f1ddb
Author: Mafi <[email protected]>
AuthorDate: Wed Sep 9 06:37:25 2026 +1000
fix(reports): bound CSV transport retries and sanitize failures (#43977)
Co-authored-by: Matt Fitzgerald <[email protected]>
Co-authored-by: Amin Ghadersohi <[email protected]>
---
docs/admin_docs/configuration/alerts-reports.mdx | 33 +++
superset/commands/report/chart_data.py | 194 ++++++++++++++
superset/commands/report/execute.py | 91 ++++---
superset/config.py | 3 +
superset/tasks/scheduler.py | 16 +-
superset/utils/csv.py | 9 +-
tests/integration_tests/reports/commands_tests.py | 109 +++++++-
.../unit_tests/commands/report/chart_data_test.py | 298 +++++++++++++++++++++
tests/unit_tests/commands/report/execute_test.py | 148 +++++++++-
.../tasks/test_scheduler_soft_timeout.py | 44 ++-
tests/unit_tests/utils/csv_tests.py | 27 +-
11 files changed, 905 insertions(+), 67 deletions(-)
diff --git a/docs/admin_docs/configuration/alerts-reports.mdx
b/docs/admin_docs/configuration/alerts-reports.mdx
index 2c873e22faa..a65dfea77f5 100644
--- a/docs/admin_docs/configuration/alerts-reports.mdx
+++ b/docs/admin_docs/configuration/alerts-reports.mdx
@@ -486,6 +486,39 @@ Log in as an admin user to ensure you have adequate
permissions.
This is the best source of information about the problem. In a docker compose
deployment, you can do this with a command like `docker logs superset_worker
--since 1h`.
+### CSV and Excel chart-data request failures
+
+The worker uses the saved query context to POST to the chart-data export
endpoint,
+falling back to the legacy GET export when a query context cannot be generated.
+`ALERT_REPORTS_CSV_REQUEST_TIMEOUT` (60 seconds by default) limits socket
operations;
+the report execution budget and its delivery/cleanup reserves also cap the
request.
+Connection and read timeouts are reported as CSV/Excel generation timeouts.
+These attachment timeouts are logged at error level and explicitly mark the
report
+task as failed, while the report execution retains its ERROR state and separate
+error-notification history. Other HTTP 408 exception handling is unchanged.
+
+To tolerate short-lived transport failures, operators can opt in with
+`ALERT_REPORTS_CSV_REQUEST_RETRY = True` (default: `False`). This permits
**one** retry
+for transient connection/read failures and HTTP 429, 500, 502, 503, or 504.
Other
+HTTP statuses are not retried. Backoff is 0.5 seconds, extended to at most 2
seconds
+for a numeric `Retry-After`; longer, invalid, or date-based delays are not
retried
+inline. Both attempts and backoff share the initial request timeout allowance
and
+respect the remaining execution budget. Unbounded requests are not retried.
+A request that consumes its entire timeout does **not** get another full
timeout.
+Socket timeouts are not wall-clock cancellation: existing report task limits
still
+interrupt in-flight work. A timed-out server query can continue running, so
enabling
+retries can increase database load. Leave retries disabled unless appropriate
for
+your deployment; disable the setting to roll back retry behavior.
+
+Worker diagnostics include schedule/chart identifiers, a fixed endpoint path
(no
+query string), error category, HTTP status, timeout, elapsed duration, and
attempt.
+For HTTP errors, at most 4097 response bytes are read to enforce a 4096-byte
limit.
+Only recognized Superset error types from up to four JSON errors are retained;
+free-form messages, extra fields, and non-JSON or oversized bodies are
redacted or
+omitted. Cookies, authentication headers, URLs, SQL, and query payloads are not
+included in these transport diagnostics. HTTP 400 therefore remains a failure
to
+investigate, not a reason to repeat the same request.
+
### Check web browser and webdriver installation
To take a screenshot, the worker visits the dashboard or chart using a
headless browser, then takes a screenshot. If you are able to send a chart as
CSV, XLSX, or text but can't send as PNG, your problem may lie with the browser.
diff --git a/superset/commands/report/chart_data.py
b/superset/commands/report/chart_data.py
new file mode 100644
index 00000000000..122c5a59c67
--- /dev/null
+++ b/superset/commands/report/chart_data.py
@@ -0,0 +1,194 @@
+# 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.
+"""Bounded, privacy-preserving transport handling for report attachments."""
+
+import errno
+import logging
+import socket
+import time
+from collections.abc import Callable
+from http.client import HTTPException, IncompleteRead
+from urllib.error import HTTPError, URLError
+
+from superset.errors import SupersetErrorType
+from superset.utils import json
+
+logger = logging.getLogger(__name__)
+_ERROR_BODY_LIMIT = 4096
+_BACKOFF_SECONDS = 0.5
+_RETRYABLE_STATUS = {429, 500, 502, 503, 504}
+_TRANSIENT_ERRNOS = {
+ errno.ECONNREFUSED,
+ errno.ECONNRESET,
+ errno.ECONNABORTED,
+ errno.ETIMEDOUT,
+ errno.EHOSTUNREACH,
+ errno.ENETUNREACH,
+ errno.EPIPE,
+}
+
+
+class ChartDataRequestError(Exception):
+ """A transport error whose message contains no remote-controlled text."""
+
+ def __init__(self, category: str, status: int | None = None) -> None:
+ """Store the category used to select the normalized report
exception."""
+ self.category = category
+ super().__init__(
+ f"Chart data request failed: category={category} status={status}"
+ )
+
+
+def _error_body(error: HTTPError) -> str:
+ """Project a bounded JSON body onto known error types, never free-form
text.
+
+ Messages, SQL, validation values and arbitrary extra fields may contain
+ credentials or query data. Redaction by keyword cannot safely retain them.
+ """
+ try:
+ body = error.read(_ERROR_BODY_LIMIT + 1)
+ if len(body) > _ERROR_BODY_LIMIT:
+ return "[omitted: oversized body]"
+ payload = json.loads(body)
+ except (ValueError, OSError, HTTPException, RecursionError):
+ return "[omitted: unreadable or non-JSON body]"
+ if not isinstance(payload, dict) or not isinstance(payload.get("errors"),
list):
+ return "[redacted]"
+ known_types = {member.value for member in SupersetErrorType}
+ errors = []
+ for item in payload["errors"][:4]:
+ error_type = item.get("error_type") if isinstance(item, dict) else None
+ errors.append(
+ {
+ "error_type": error_type
+ if isinstance(error_type, str) and error_type in known_types
+ else "[redacted]",
+ "message": "[redacted]",
+ }
+ )
+ return json.dumps({"errors": errors})
+
+
+def _retry_delay(error: HTTPError) -> float | None:
+ """Respect numeric Retry-After; defer date-based delays to the
scheduler."""
+ retry_after = error.headers.get("Retry-After") if error.headers else None
+ if retry_after is None:
+ return _BACKOFF_SECONDS
+ try:
+ seconds = float(retry_after)
+ except ValueError:
+ return None
+ return max(_BACKOFF_SECONDS, seconds) if 0 <= seconds <= 2 else None
+
+
+def request_chart_data(
+ fetch: Callable[[float | None], bytes | None],
+ get_timeout: Callable[[], float | None],
+ *,
+ retry: bool,
+ endpoint: str,
+ log_context: str,
+) -> bytes | None:
+ """Optionally retry once, sharing the initial timeout and execution budget.
+
+ The caller supplies a fixed endpoint path, never a URL or query payload.
+ The deadline callback also preserves the report's delivery/cleanup
reserves.
+ Socket timeouts are not wall-clock cancellation; the report task's existing
+ execution limits remain responsible for interrupting in-flight work.
+ """
+ started = time.monotonic()
+ timeout = get_timeout()
+ deadline = started + timeout if timeout is not None else None
+ for attempt in (1, 2):
+ if attempt == 2:
+ phase_timeout = get_timeout()
+ remaining = deadline - time.monotonic() if deadline is not None
else 0
+ if remaining <= 0:
+ raise ChartDataRequestError("timeout")
+ timeout = (
+ min(remaining, phase_timeout)
+ if phase_timeout is not None
+ else remaining
+ )
+ status = None
+ body = "[not applicable]"
+ delay = _BACKOFF_SECONDS
+ try:
+ content = fetch(timeout)
+ logger.info(
+ "Chart data request completed %s endpoint=%s "
+ "elapsed_seconds=%.2f attempt=%s timeout_seconds=%s",
+ log_context,
+ endpoint,
+ time.monotonic() - started,
+ attempt,
+ timeout,
+ )
+ return content
+ except HTTPError as error:
+ status = error.code
+ category = "http"
+ transient = status in _RETRYABLE_STATUS
+ # A server requesting a longer or date-based delay should be left
+ # to the scheduler, rather than retried earlier than requested.
+ retry_delay = _retry_delay(error)
+ if retry_delay is None:
+ transient = False
+ else:
+ delay = retry_delay
+ try:
+ body = _error_body(error)
+ finally:
+ error.close()
+ except (OSError, HTTPException) as error:
+ reason = error.reason if isinstance(error, URLError) else error
+ is_timeout = isinstance(reason, TimeoutError) or (
+ isinstance(reason, OSError) and reason.errno == errno.ETIMEDOUT
+ )
+ category = "timeout" if is_timeout else "network"
+ transient = (
+ is_timeout
+ or isinstance(reason, IncompleteRead)
+ or (isinstance(reason, OSError) and reason.errno in
_TRANSIENT_ERRNOS)
+ or (
+ isinstance(reason, socket.gaierror)
+ and reason.errno == socket.EAI_AGAIN
+ )
+ )
+ logger.warning(
+ "Chart data request failed %s endpoint=%s category=%s status=%s "
+ "elapsed_seconds=%.2f attempt=%s timeout_seconds=%s body=%s",
+ log_context,
+ endpoint,
+ category,
+ status,
+ time.monotonic() - started,
+ attempt,
+ timeout,
+ body,
+ )
+ # Never retry an unbounded request, or renew the original timeout.
+ if not retry or not transient or attempt == 2 or deadline is None:
+ raise ChartDataRequestError(category, status) from None
+ phase_timeout = get_timeout()
+ remaining = deadline - time.monotonic()
+ if phase_timeout is not None:
+ remaining = min(remaining, phase_timeout)
+ if remaining <= delay:
+ raise ChartDataRequestError(category, status) from None
+ time.sleep(delay)
+ return None # pragma: no cover
diff --git a/superset/commands/report/execute.py
b/superset/commands/report/execute.py
index b385db05e1c..17548c26efe 100644
--- a/superset/commands/report/execute.py
+++ b/superset/commands/report/execute.py
@@ -35,6 +35,10 @@ from superset.commands.base import BaseCommand
from superset.commands.dashboard.permalink.create import
CreateDashboardPermalinkCommand
from superset.commands.exceptions import CommandException, UpdateFailedError
from superset.commands.report.alert import AlertCommand
+from superset.commands.report.chart_data import (
+ ChartDataRequestError,
+ request_chart_data,
+)
from superset.commands.report.exceptions import (
ReportScheduleAlertGracePeriodError,
ReportScheduleClientErrorsException,
@@ -997,7 +1001,7 @@ class BaseReportState:
raise URLError(response.getcode())
return content or None
- def _get_data(self, result_format: ChartDataResultFormat) -> bytes:
+ def _get_data(self, result_format: ChartDataResultFormat) -> bytes: #
noqa: C901
"""
Fetch tabular chart data (CSV or Excel) as raw bytes.
@@ -1035,50 +1039,56 @@ class BaseReportState:
self._update_query_context(failed_error)
db.session.refresh(self._report_schedule.chart)
+ def get_timeout() -> float | None:
+ """Cap every request by the available data-generation budget."""
+ return self._phase_timeout(
+ "data_generation",
+
requested_seconds=app.config["ALERT_REPORTS_CSV_REQUEST_TIMEOUT"],
+ reserve_seconds=(
+ self._report_execution_context.post_capture_reserve_seconds
+ if self._report_execution_context
+ else 0.0
+ ),
+ )
+
try:
if self._report_schedule.chart.query_context is None:
url = self._get_url(result_format=result_format)
- data = get_chart_csv_data(
- chart_url=url,
- auth_cookies=auth_cookies,
- timeout=self._phase_timeout(
- "data_generation",
- requested_seconds=app.config[
- "ALERT_REPORTS_CSV_REQUEST_TIMEOUT"
- ],
- reserve_seconds=(
-
self._report_execution_context.post_capture_reserve_seconds
- if self._report_execution_context
- else 0.0
- ),
- ),
- )
+ endpoint = "/api/v1/chart/{id}/data/"
+
+ def fetch(timeout: float | None) -> bytes | None:
+ """Fetch the legacy export without exposing its URL in
logs."""
+ return get_chart_csv_data(
+ chart_url=url, auth_cookies=auth_cookies,
timeout=timeout
+ )
else:
request_payload =
self._get_chart_data_request_payload(result_format)
url = get_url_path("ChartDataRestApi.data")
- data = self._post_chart_data(
- chart_url=url,
- auth_cookies=auth_cookies,
- request_payload=request_payload,
- timeout=self._phase_timeout(
- "data_generation",
- requested_seconds=app.config[
- "ALERT_REPORTS_CSV_REQUEST_TIMEOUT"
- ],
- reserve_seconds=(
-
self._report_execution_context.post_capture_reserve_seconds
- if self._report_execution_context
- else 0.0
- ),
- ),
- )
+ endpoint = "/api/v1/chart/data"
+
+ def fetch(timeout: float | None) -> bytes | None:
+ """Use the saved query context's existing POST export
path."""
+ return self._post_chart_data(
+ chart_url=url,
+ auth_cookies=auth_cookies,
+ request_payload=request_payload,
+ timeout=timeout,
+ )
+
+ data = request_chart_data(
+ fetch,
+ get_timeout,
+ retry=app.config["ALERT_REPORTS_CSV_REQUEST_RETRY"],
+ endpoint=endpoint,
+ log_context=self._log_context,
+ )
elapsed_seconds: float = (
datetime.now(timezone.utc).replace(tzinfo=None) - start_time
).total_seconds()
logger.info(
"%s data generation from %s as user %s took %.2fs -
execution_id: %s",
label,
- url,
+ endpoint,
username,
elapsed_seconds,
self._execution_id,
@@ -1096,6 +1106,10 @@ class BaseReportState:
if self._report_schedule.type == ReportScheduleType.REPORT:
raise
raise timeout_error() from ex
+ except ChartDataRequestError as ex:
+ if ex.category == "timeout":
+ raise timeout_error() from ex
+ raise failed_error(str(ex)) from ex
except ReportExecutionBudgetExceededError:
raise
except Exception as ex:
@@ -1845,9 +1859,9 @@ class ReportNotTriggeredErrorState(BaseReportState):
second_error_message = str(second_ex)
finally:
try:
- self.update_report_schedule_and_log(
- ReportState.ERROR,
- error_message=second_error_message,
+ # Notification bookkeeping is not another execution
outcome.
+ self.create_log(
+ second_error_message,
include_execution_warnings=False,
)
except ReportScheduleUnexpectedError:
@@ -2050,8 +2064,9 @@ class ReportSuccessState(BaseReportState):
second_error_message = str(second_ex)
finally:
try:
- self.update_report_schedule_and_log(
- ReportState.ERROR,
error_message=second_error_message
+ # Preserve the grace-period marker without another
terminal log.
+ self.create_log(
+ second_error_message,
include_execution_warnings=False
)
except ReportScheduleUnexpectedError:
# Logging failed again; log it but don't hide first_ex
diff --git a/superset/config.py b/superset/config.py
index 00091152314..3a32ae85b81 100644
--- a/superset/config.py
+++ b/superset/config.py
@@ -2545,6 +2545,9 @@ ALERT_REPORTS_QUERY_EXECUTION_MAX_TRIES = 1
# which leaves the report schedule stuck in the WORKING state. Set to None to
# disable (not recommended).
ALERT_REPORTS_CSV_REQUEST_TIMEOUT = 60
+# Opt in to at most one transient CSV/Excel transport retry within the original
+# request timeout and report execution budget. Does not retry unbounded
requests.
+ALERT_REPORTS_CSV_REQUEST_RETRY = False
# Custom width for screenshots
ALERT_REPORTS_MIN_CUSTOM_SCREENSHOT_WIDTH = 600
ALERT_REPORTS_MAX_CUSTOM_SCREENSHOT_WIDTH = 2400
diff --git a/superset/tasks/scheduler.py b/superset/tasks/scheduler.py
index 8581ac424b0..0e1fef705b4 100644
--- a/superset/tasks/scheduler.py
+++ b/superset/tasks/scheduler.py
@@ -31,7 +31,11 @@ from superset_core.tasks.types import TaskStatus
from superset import is_feature_enabled
from superset.commands.exceptions import CommandException
from superset.commands.logs.prune import LogPruneCommand
-from superset.commands.report.exceptions import ReportScheduleUnexpectedError
+from superset.commands.report.exceptions import (
+ ReportScheduleCsvTimeout,
+ ReportScheduleUnexpectedError,
+ ReportScheduleXlsxTimeout,
+)
from superset.commands.report.execute import AsyncExecuteReportScheduleCommand
from superset.commands.report.log_prune import
AsyncPruneReportScheduleLogCommand
from superset.commands.sql_lab.query import QueryPruneCommand
@@ -178,6 +182,16 @@ def execute(
"An unexpected error occurred while executing the report: %s",
task_id
)
self.update_state(state="FAILURE")
+ except (ReportScheduleCsvTimeout, ReportScheduleXlsxTimeout):
+ # Attachment generation timeouts are failed executions, despite their
+ # HTTP 408 status. Keep them visible to error-level task monitoring.
+ logger.exception(
+ "Report attachment generation timed out; execution_id=%s "
+ "report_schedule_id=%s",
+ task_id,
+ report_schedule_id,
+ )
+ self.update_state(state="FAILURE")
except CommandException as ex:
logger_func, level = get_logger_from_status(ex.status)
logger_func(
diff --git a/superset/utils/csv.py b/superset/utils/csv.py
index 20d17702c1d..c0c9b10dcc8 100644
--- a/superset/utils/csv.py
+++ b/superset/utils/csv.py
@@ -16,6 +16,7 @@
# under the License.
import logging
import urllib.request
+from contextlib import closing
from typing import Any, Optional, Union
from urllib.error import URLError
@@ -119,10 +120,10 @@ def get_chart_csv_data(
opener.addheaders.append(("Cookie", cookie_str))
# A missing timeout means the socket blocks forever when the Superset
# webserver is unreachable, wedging the report schedule in WORKING.
- response = opener.open(chart_url, timeout=timeout)
- content = response.read()
- if response.getcode() != 200:
- raise URLError(response.getcode())
+ with closing(opener.open(chart_url, timeout=timeout)) as response:
+ content = response.read()
+ if response.getcode() != 200:
+ raise URLError(response.getcode())
if content:
return content
return None
diff --git a/tests/integration_tests/reports/commands_tests.py
b/tests/integration_tests/reports/commands_tests.py
index c4aa76a9147..01654e6f2a5 100644
--- a/tests/integration_tests/reports/commands_tests.py
+++ b/tests/integration_tests/reports/commands_tests.py
@@ -2878,20 +2878,29 @@ def
test_readiness_timeout_retries_terminal_persistence_and_allows_next_schedule
"load_birth_names_dashboard_with_slices",
"create_report_email_chart_with_csv"
)
@patch("superset.reports.notifications.email.send_email_smtp")
-@patch("superset.utils.csv.urllib.request.urlopen")
@patch("superset.utils.csv.urllib.request.OpenerDirector.open")
-@patch("superset.utils.csv.get_chart_csv_data")
def test_fail_csv(
- csv_mock, mock_open, mock_urlopen, email_mock,
create_report_email_chart_with_csv
-):
+ mock_open: Mock,
+ email_mock: Mock,
+ create_report_email_chart_with_csv: ReportSchedule,
+ caplog: pytest.LogCaptureFixture,
+) -> None:
"""
ExecuteReport Command: Test error on csv
"""
- response = Mock()
- mock_open.return_value = response
- mock_urlopen.return_value = response
- mock_urlopen.return_value.getcode.return_value = 500
+ from email.message import Message
+ from io import BytesIO
+ from urllib.error import HTTPError
+
+ caplog.set_level(logging.INFO, logger="superset.commands.report.execute")
+ mock_open.side_effect = HTTPError(
+ "http://localhost/api/v1/chart/data",
+ 500,
+ "Internal Server Error",
+ Message(),
+ BytesIO(b'{"message":"error details"}'),
+ )
with pytest.raises(ReportScheduleCsvFailedError):
AsyncExecuteReportScheduleCommand(
@@ -2903,9 +2912,18 @@ def test_fail_csv(
assert email_mock.call_args[0][0] == DEFAULT_OWNER_EMAIL
assert_log(
- ReportState.ERROR, error_message="Failed generating csv <urlopen error
500>"
+ ReportState.ERROR,
+ error_message="Chart data request failed: category=http status=500",
)
+ terminals = [
+ record.message
+ for record in caplog.records
+ if "report_execution_terminal" in record.message
+ ]
+ assert len(terminals) == 1
+ assert "category=http status=500" in terminals[0]
+
@pytest.mark.usefixtures(
"load_birth_names_dashboard_with_slices", "create_alert_email_chart"
@@ -3568,3 +3586,76 @@ def test_get_retry_delay_exponential_backoff() -> None:
assert state._get_retry_delay(5) == 1920 # 60 * 2^5 = 1920
assert state._get_retry_delay(6) == 3600 # 60 * 2^6 = 3840 → capped at
3600
assert state._get_retry_delay(10) == 3600 # still capped
+
+
[email protected]("load_birth_names_dashboard_with_slices")
[email protected]("report_format", [ReportDataFormat.CSV,
ReportDataFormat.XLSX])
[email protected]("post", [False, True])
[email protected]("wrapped", [False, True])
[email protected]("retry", [False, True])
+def test_scheduler_tabular_transport_timeout_records_failure(
+ create_report_email_chart_with_csv: ReportSchedule,
+ report_format: ReportDataFormat,
+ post: bool,
+ wrapped: bool,
+ retry: bool,
+ caplog: pytest.LogCaptureFixture,
+) -> None:
+ """Socket timeouts remain task failures with one ERROR outcome and
notification."""
+ import socket
+ from urllib.error import URLError
+
+ from superset.commands.report.exceptions import (
+ ReportScheduleCsvTimeout,
+ ReportScheduleXlsxTimeout,
+ )
+ from superset.tasks.scheduler import execute
+
+ schedule = create_report_email_chart_with_csv
+ schedule.report_format = report_format
+ schedule.chart.query_context = "{}" if post else None
+ db.session.commit()
+ timeout = socket.timeout("TRANSPORT_SECRET")
+ error = URLError(timeout) if wrapped else timeout
+ expected = (
+ ReportScheduleCsvTimeout
+ if report_format == ReportDataFormat.CSV
+ else ReportScheduleXlsxTimeout
+ )()
+ caplog.set_level(logging.INFO)
+ with (
+ patch.dict(app.config, {"ALERT_REPORTS_CSV_REQUEST_RETRY": retry}),
+ patch.object(BaseReportState, "_update_query_context"),
+ patch(
+ "superset.utils.csv.urllib.request.OpenerDirector.open",
side_effect=error
+ ) as fetch,
+ patch("superset.commands.report.chart_data.time.sleep"),
+ patch("superset.reports.notifications.email.send_email_smtp") as email,
+ patch.object(execute, "update_state") as update_state,
+ ):
+ execute.push_request(id=TEST_ID)
+ try:
+ execute.run(schedule.id)
+ finally:
+ execute.pop_request()
+
+ assert fetch.call_count == (2 if retry else 1)
+ assert isinstance(fetch.call_args.args[0], str) is not post
+ update_state.assert_called_once_with(state="FAILURE")
+ errors = [
+ record
+ for record in caplog.records
+ if record.name == "superset.tasks.scheduler" and record.levelno >=
logging.ERROR
+ ]
+ assert len(errors) == 1
+ assert errors[0].exc_info is not None
+ assert isinstance(errors[0].exc_info[1], type(expected))
+ assert "TRANSPORT_SECRET" not in caplog.text
+ db.session.refresh(schedule)
+ assert schedule.last_state == ReportState.ERROR
+ assert_log(ReportState.ERROR, error_message=str(expected))
+ assert (
+ sum("report_execution_terminal" in record.message for record in
caplog.records)
+ == 1
+ )
+ email.assert_called_once()
diff --git a/tests/unit_tests/commands/report/chart_data_test.py
b/tests/unit_tests/commands/report/chart_data_test.py
new file mode 100644
index 00000000000..3dc2bafb193
--- /dev/null
+++ b/tests/unit_tests/commands/report/chart_data_test.py
@@ -0,0 +1,298 @@
+# 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.
+"""Report transport retries must be bounded and must not disclose request
data."""
+
+import errno
+import io
+import socket
+from email.message import Message
+from http.client import IncompleteRead
+from unittest.mock import Mock
+from urllib.error import HTTPError, URLError
+
+import pytest
+from pytest_mock import MockerFixture
+
+from superset.commands.report.chart_data import (
+ ChartDataRequestError,
+ request_chart_data,
+)
+from superset.utils import json
+from superset.utils.report_execution import (
+ ReportExecutionBudgetExceededError,
+ ReportExecutionDeadline,
+)
+
+
[email protected]
+def clock(mocker: MockerFixture) -> list[float]:
+ """Advance a deterministic monotonic clock when the retry sleeps."""
+ now = [0.0]
+ mocker.patch(
+ "superset.commands.report.chart_data.time.monotonic",
side_effect=lambda: now[0]
+ )
+ mocker.patch(
+ "superset.commands.report.chart_data.time.sleep",
+ side_effect=lambda delay: now.__setitem__(0, now[0] + delay),
+ )
+ return now
+
+
+def request(
+ fetch: Mock, timeout: float | None = 10, retry: bool = True
+) -> bytes | None:
+ """Invoke the helper using only safe, static diagnostic context."""
+ return request_chart_data(
+ fetch,
+ lambda: timeout,
+ retry=retry,
+ endpoint="/api/v1/chart/data",
+ log_context="report_schedule_id=1 chart_id=2",
+ )
+
+
+def http_error(
+ status: int, body: bytes = b"", retry_after: str | None = None
+) -> HTTPError:
+ """Build an error with sensitive URL, reason and headers that must not
leak."""
+ headers = Message()
+ headers["Set-Cookie"] = "session=HEADER_SECRET"
+ if retry_after is not None:
+ headers["Retry-After"] = retry_after
+ return HTTPError(
+ "https://user:URL_SECRET@localhost/path?form_data=QUERY_SECRET",
+ status,
+ "REASON_SECRET",
+ headers,
+ io.BytesIO(body),
+ )
+
+
[email protected]("status", [400, 401, 403, 404, 422, 501, 505])
+def test_http_deterministic_error_is_sanitized(
+ status: int,
+ clock: list[float],
+ caplog: pytest.LogCaptureFixture,
+) -> None:
+ """Preserve known error types but no arbitrary messages, extras or
credentials."""
+ body = json.dumps(
+ {
+ "errors": [
+ {
+ "error_type": "INVALID_PAYLOAD_SCHEMA_ERROR",
+ "message": "MESSAGE_SECRET",
+ "extra": {"sql": "PAYLOAD_SECRET"},
+ },
+ {"error_type": "TYPE_SECRET", "message": "SECOND_SECRET"},
+ ],
+ "token": "TOKEN_SECRET",
+ }
+ ).encode()
+ error = http_error(status, body)
+ fetch = Mock(side_effect=error)
+ with pytest.raises(
+ ChartDataRequestError, match=f"category=http status={status}"
+ ) as exc:
+ request(fetch)
+ fetch.assert_called_once_with(10)
+ assert clock == [0]
+ assert error.closed
+ output = caplog.text + str(exc.value)
+ assert "SECRET" not in output
+ assert "INVALID_PAYLOAD_SCHEMA_ERROR" in output
+ assert "[redacted]" in output
+ assert "report_schedule_id=1 chart_id=2" in output
+ assert "endpoint=/api/v1/chart/data" in output
+ assert "attempt=1" in output
+ assert "elapsed_seconds=0.00" in output
+ assert len(caplog.messages[0]) < 1024
+ assert exc.value.__cause__ is None
+
+
[email protected](
+ "body",
+ [b"x" * 10000, b"<html>BODY_SECRET</html>", b'{"message":"BODY_SECRET"}',
b"\xff"],
+)
+def test_error_body_is_bounded(
+ body: bytes,
+ clock: list[float],
+ caplog: pytest.LogCaptureFixture,
+ mocker: MockerFixture,
+) -> None:
+ """Read at most the fixed limit plus one sentinel byte and log no raw
text."""
+ error = http_error(400, body)
+ read = mocker.patch.object(error, "read", wraps=error.read)
+ with pytest.raises(ChartDataRequestError):
+ request(Mock(side_effect=error))
+ read.assert_called_once_with(4097)
+ assert "SECRET" not in caplog.text
+ assert len(caplog.messages[0]) < 1024
+
+
[email protected]("status", [429, 500, 502, 503, 504])
+def test_retryable_http_success(status: int, clock: list[float]) -> None:
+ """Transient HTTP responses may retry once using the remaining timeout."""
+ fetch = Mock(side_effect=[http_error(status), b"csv"])
+ assert request(fetch) == b"csv"
+ assert [call.args[0] for call in fetch.call_args_list] == [10, 9.5]
+ assert clock == [0.5]
+
+
[email protected](
+ "error",
+ [
+ IncompleteRead(b"PAYLOAD_SECRET"),
+ socket.timeout("TIMEOUT_SECRET"),
+ URLError(TimeoutError("TIMEOUT_SECRET")),
+ URLError(ConnectionRefusedError(errno.ECONNREFUSED, "NETWORK_SECRET")),
+ ConnectionResetError(errno.ECONNRESET, "NETWORK_SECRET"),
+ URLError(socket.gaierror(socket.EAI_AGAIN, "DNS_SECRET")),
+ ],
+)
+def test_transient_network_retry(
+ error: Exception, clock: list[float], caplog: pytest.LogCaptureFixture
+) -> None:
+ """Direct read errors and wrapped connection errors share the bounded
retry."""
+ fetch = Mock(side_effect=[error, b"csv"])
+ assert request(fetch) == b"csv"
+ assert fetch.call_count == 2
+ assert "SECRET" not in caplog.text
+
+
[email protected](
+ "error,category",
+ [(TimeoutError("TIMEOUT_SECRET"), "timeout"), (http_error(503), "http")],
+)
+def test_retry_exhaustion(error: OSError, category: str, clock: list[float])
-> None:
+ """An opt-in retry never becomes an unbounded loop."""
+ # Each HTTP failure needs its own response stream.
+ fetch = Mock(side_effect=[error, http_error(503) if category == "http"
else error])
+ with pytest.raises(ChartDataRequestError, match=f"category={category}"):
+ request(fetch)
+ assert fetch.call_count == 2
+ assert clock == [0.5]
+
+
[email protected]("timeout,retry", [(10, False), (None, True), (0.5,
True)])
+def test_retry_disabled_or_no_budget(
+ timeout: float | None, retry: bool, clock: list[float]
+) -> None:
+ """Defaults, unbounded requests and insufficient budgets do not retry."""
+ fetch = Mock(side_effect=TimeoutError())
+ with pytest.raises(ChartDataRequestError, match="category=timeout"):
+ request(fetch, timeout, retry)
+ fetch.assert_called_once()
+ assert clock == [0]
+
+
+def test_read_timeout_consumes_original_budget(clock: list[float]) -> None:
+ """A full-length read timeout does not get another full request budget."""
+
+ def fail(timeout: float | None) -> bytes:
+ """Simulate a request using all its allowed time."""
+ clock[0] += 10
+ raise TimeoutError()
+
+ fetch = Mock(side_effect=fail)
+ with pytest.raises(ChartDataRequestError, match="category=timeout"):
+ request(fetch)
+ fetch.assert_called_once()
+ assert clock == [10]
+
+
+def test_report_deadline_prevents_retry(clock: list[float]) -> None:
+ """Respect the execution deadline and preserve its phase reserves."""
+ deadline = ReportExecutionDeadline(
+ total_seconds=10, started_at=0, _clock=lambda: clock[0]
+ )
+
+ def fail(timeout: float | None) -> bytes:
+ """Consume the work allowance without consuming the cleanup reserve."""
+ clock[0] = 8
+ raise TimeoutError()
+
+ fetch = Mock(side_effect=fail)
+ with pytest.raises(ReportExecutionBudgetExceededError):
+ request_chart_data(
+ fetch,
+ lambda: deadline.timeout_seconds(
+ "data_generation", requested_seconds=60, reserve_seconds=2
+ ),
+ retry=True,
+ endpoint="/api/v1/chart/data",
+ log_context="chart_id=2",
+ )
+ fetch.assert_called_once_with(8)
+ assert clock == [8]
+
+
[email protected](
+ "retry_after, calls",
+ [("2", 2), ("20", 1), ("Wed, 01 Jan 2030 00:00:00 GMT", 1), ("nan", 1)],
+)
+def test_retry_after(retry_after: str, calls: int, clock: list[float]) -> None:
+ """Never retry earlier than Retry-After or sleep beyond the budget."""
+ fetch = Mock(side_effect=[http_error(429, retry_after=retry_after),
b"csv"])
+ if calls == 1:
+ with pytest.raises(ChartDataRequestError):
+ request(fetch)
+ assert clock == [0]
+ else:
+ assert request(fetch) == b"csv"
+ assert clock == [2]
+ assert fetch.call_count == calls
+
+
+def test_non_transient_url_error(
+ clock: list[float], caplog: pytest.LogCaptureFixture
+) -> None:
+ """String reasons, including TLS errors, must not be blindly retried or
logged."""
+ fetch = Mock(side_effect=URLError("URL_SECRET"))
+ with pytest.raises(ChartDataRequestError, match="category=network"):
+ request(fetch)
+ fetch.assert_called_once()
+ assert "SECRET" not in caplog.text
+
+
[email protected]("read_error", [TimeoutError(),
IncompleteRead(b"BODY_SECRET")])
+def test_http_error_body_read_failure(
+ read_error: Exception,
+ clock: list[float],
+ caplog: pytest.LogCaptureFixture,
+ mocker: MockerFixture,
+) -> None:
+ """An unreadable diagnostic body must not replace the original HTTP
failure."""
+ error = http_error(400)
+ mocker.patch.object(error, "read", side_effect=read_error)
+ with pytest.raises(ChartDataRequestError, match="category=http
status=400"):
+ request(Mock(side_effect=error))
+ assert error.closed
+ assert "SECRET" not in caplog.text
+
+
+def test_deadline_is_rechecked_after_backoff(
+ clock: list[float], mocker: MockerFixture
+) -> None:
+ """An oversleep cannot start another request after the shared budget
expires."""
+ mocker.patch(
+ "superset.commands.report.chart_data.time.sleep",
+ side_effect=lambda _: clock.__setitem__(0, 11),
+ )
+ fetch = Mock(side_effect=TimeoutError())
+ with pytest.raises(ChartDataRequestError, match="category=timeout"):
+ request(fetch)
+ fetch.assert_called_once()
diff --git a/tests/unit_tests/commands/report/execute_test.py
b/tests/unit_tests/commands/report/execute_test.py
index d6f6a7dee6f..5bf2da33dd8 100644
--- a/tests/unit_tests/commands/report/execute_test.py
+++ b/tests/unit_tests/commands/report/execute_test.py
@@ -3749,6 +3749,7 @@ def test_success_state_send_error_logs_and_reraises(
)
mocker.patch.object(state, "send", side_effect=RuntimeError("send boom"))
mocker.patch.object(state, "is_in_error_grace_period", return_value=False)
+ mocker.patch.object(state, "create_log")
mocker.patch.object(state, "send_error")
mocker.patch.object(state, "update_report_schedule_and_log")
@@ -4280,6 +4281,7 @@ def test_success_state_send_failure_notifies_owner(
mocker.patch.object(state, "is_in_grace_period", return_value=False)
mocker.patch.object(state, "is_in_error_grace_period", return_value=False)
mock_update = mocker.patch.object(state, "update_report_schedule_and_log")
+ mocker.patch.object(state, "create_log")
mock_send_error = mocker.patch.object(state, "send_error")
if schedule_type == ReportScheduleType.ALERT:
mocker.patch(
@@ -4339,10 +4341,11 @@ def
test_success_state_send_error_failure_overwrites_marker(
expected_message: str,
) -> None:
"""When the Success/Grace path's own error notification fails, the
- placeholder marker is overwritten with the real failure message before
- ERROR is logged -- mirroring the first-run (ReportNotTriggeredErrorState)
- path. A SupersetErrorsException contributes its joined error messages; any
- other exception contributes its ``str()``."""
+ notification audit entry records the real failure message instead of the
+ success marker, without replacing the terminal execution outcome. This
mirrors
+ the first-run (ReportNotTriggeredErrorState) path. A
SupersetErrorsException
+ contributes its joined error messages; any other exception contributes its
+ ``str()``."""
from superset.errors import ErrorLevel, SupersetError, SupersetErrorType
from superset.exceptions import SupersetErrorsException
@@ -4369,6 +4372,7 @@ def
test_success_state_send_error_failure_overwrites_marker(
)
mocker.patch.object(state, "is_in_error_grace_period", return_value=False)
mock_update = mocker.patch.object(state, "update_report_schedule_and_log")
+ mock_log = mocker.patch.object(state, "create_log")
mock_send_error = mocker.patch.object(
state, "send_error", side_effect=send_error_exc
)
@@ -4382,15 +4386,9 @@ def
test_success_state_send_error_failure_overwrites_marker(
state.next()
mock_send_error.assert_called_once()
- # The placeholder marker must be replaced by the real notification failure
- # before the terminal ERROR row is written.
- final_call = mock_update.call_args_list[-1]
- assert final_call.args[0] == ReportState.ERROR
- assert final_call.kwargs.get("error_message") == expected_message
- assert (
- final_call.kwargs.get("error_message")
- != REPORT_SCHEDULE_ERROR_NOTIFICATION_MARKER
- )
+ # Notification bookkeeping preserves the original execution failure.
+ assert mock_update.call_args_list[-1].kwargs["error_message"] == "blank
screenshot"
+ mock_log.assert_called_once_with(expected_message,
include_execution_warnings=False)
def test_get_url_for_csv_uses_post_processed_type(
@@ -4562,3 +4560,127 @@ def
test_get_url_raises_unexpected_error_when_target_is_missing(
assert "orphan_report" in message
assert "chart_id=None" in message
assert "dashboard_id=None" in message
+
+
[email protected](
+ "state_class", [ReportNotTriggeredErrorState, ReportSuccessState]
+)
[email protected]("notification_fails", [False, True])
+def test_report_failure_has_one_terminal_outcome(
+ mocker: MockerFixture,
+ caplog: pytest.LogCaptureFixture,
+ state_class: type[BaseReportState],
+ notification_fails: bool,
+) -> None:
+ """Keep notification bookkeeping separate from the execution's terminal
log."""
+ state = _make_state_instance(
+ mocker, state_class, schedule_type=ReportScheduleType.REPORT
+ )
+ mocker.patch.object(
+ state, "send", side_effect=ReportScheduleCsvFailedError("export
failed")
+ )
+ mocker.patch.object(state, "is_in_error_grace_period", return_value=False)
+ mock_log = mocker.patch.object(state, "create_log")
+ send_error = mocker.patch.object(
+ state,
+ "send_error",
+ side_effect=RuntimeError("notification failed") if notification_fails
else None,
+ )
+ caplog.set_level("INFO", logger="superset.commands.report.execute")
+ with pytest.raises(ReportScheduleCsvFailedError, match="export failed"):
+ state.next()
+ terminals = [
+ record.message
+ for record in caplog.records
+ if "report_execution_terminal" in record.message
+ ]
+ assert len(terminals) == 1
+ assert "export failed" in terminals[0]
+ assert state._report_schedule.last_state == ReportState.ERROR
+ send_error.assert_called_once()
+ assert mock_log.call_args.args[0] == (
+ "notification failed"
+ if notification_fails
+ else REPORT_SCHEDULE_ERROR_NOTIFICATION_MARKER
+ )
+ assert mock_log.call_args.kwargs == {"include_execution_warnings": False}
+
+
[email protected]("post", [False, True])
[email protected]("wrapped", [False, True])
[email protected](
+ "result_format", [ChartDataResultFormat.CSV, ChartDataResultFormat.XLSX]
+)
+def test_chart_data_normalizes_transport_timeouts(
+ app: SupersetApp,
+ mocker: MockerFixture,
+ caplog: pytest.LogCaptureFixture,
+ post: bool,
+ wrapped: bool,
+ result_format: ChartDataResultFormat,
+) -> None:
+ """Both export paths normalize direct and urllib-wrapped socket
timeouts."""
+ from superset.commands.report.exceptions import (
+ ReportScheduleCsvTimeout,
+ ReportScheduleXlsxTimeout,
+ )
+
+ state = BaseReportState(create_report_schedule(mocker), datetime.utcnow(),
uuid4())
+ _mock_xlsx_chart_data_dependencies(mocker, state)
+ if post:
+ state._report_schedule.chart.query_context = "{}"
+ error = TimeoutError("SECRET timeout reason")
+ fetch = mocker.patch(
+ "superset.commands.report.execute.BaseReportState._post_chart_data"
+ if post
+ else "superset.commands.report.execute.get_chart_csv_data",
+ side_effect=URLError(error) if wrapped else error,
+ )
+ expected = (
+ ReportScheduleCsvTimeout
+ if result_format == ChartDataResultFormat.CSV
+ else ReportScheduleXlsxTimeout
+ )
+ with pytest.raises(expected) as exc:
+ state._get_data(result_format)
+ fetch.assert_called_once()
+ assert "SECRET" not in caplog.text + str(exc.value)
+
+
[email protected]("post", [False, True])
+def test_chart_data_http_failure_does_not_expose_request(
+ app: SupersetApp,
+ mocker: MockerFixture,
+ caplog: pytest.LogCaptureFixture,
+ post: bool,
+) -> None:
+ """Transport failure diagnostics never include cookies or export
payloads."""
+ import traceback
+ from email.message import Message
+ from io import BytesIO
+ from urllib.error import HTTPError
+
+ state = BaseReportState(create_report_schedule(mocker), datetime.utcnow(),
uuid4())
+ get_url, cookies = _mock_xlsx_chart_data_dependencies(mocker, state)
+ cookies["session"] = "COOKIE_SECRET"
+ get_url.return_value = (
+ "https://localhost/api/v1/chart/1/data?form_data=QUERY_SECRET"
+ )
+ if post:
+ state._report_schedule.chart.query_context = (
+ '{"queries":[{"sql":"PAYLOAD_SECRET"}]}'
+ )
+ mocker.patch.dict(app.config, {"ALERT_REPORTS_CSV_REQUEST_RETRY": True})
+ opener = mocker.patch("urllib.request.build_opener").return_value
+ opener.open.side_effect = HTTPError(
+ "https://localhost/chart?form_data=QUERY_SECRET",
+ 400,
+ "REASON_SECRET",
+ Message(),
+ BytesIO(b'{"message":"BODY_SECRET"}'),
+ )
+ with pytest.raises(ReportScheduleCsvFailedError, match="status=400") as
exc:
+ state._get_data(ChartDataResultFormat.CSV)
+ opener.open.assert_called_once()
+ assert "SECRET" not in caplog.text + str(exc.value)
+ assert "SECRET" not in "".join(traceback.format_exception(exc.value))
diff --git a/tests/unit_tests/tasks/test_scheduler_soft_timeout.py
b/tests/unit_tests/tasks/test_scheduler_soft_timeout.py
index 6eb0a80e776..2c67dfb6764 100644
--- a/tests/unit_tests/tasks/test_scheduler_soft_timeout.py
+++ b/tests/unit_tests/tasks/test_scheduler_soft_timeout.py
@@ -14,7 +14,7 @@
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
-"""Unit tests for the shared ``reports.execute`` soft-timeout handler."""
+"""Unit tests for report task timeout handling and failure monitoring."""
from unittest.mock import MagicMock, patch
@@ -66,3 +66,45 @@ def test_soft_timeout_handler_is_shared_by_alerts() -> None:
and alert_schedule_id in call.args
for call in logger_mock.warning.call_args_list
)
+
+
[email protected]("attachment", ["csv", "xlsx", "screenshot"])
+def test_attachment_timeout_failure_policy(
+ attachment: str,
+ caplog: pytest.LogCaptureFixture,
+) -> None:
+ """Tabular generation timeouts are failures without changing other 408
policy."""
+ import logging
+
+ from superset.commands.report.exceptions import (
+ ReportScheduleCsvTimeout,
+ ReportScheduleScreenshotTimeout,
+ ReportScheduleXlsxTimeout,
+ )
+ from superset.tasks.scheduler import execute
+
+ error = {
+ "csv": ReportScheduleCsvTimeout,
+ "xlsx": ReportScheduleXlsxTimeout,
+ "screenshot": ReportScheduleScreenshotTimeout,
+ }[attachment]()
+ assert error.status == 408
+ caplog.set_level(logging.WARNING)
+ with (
+ patch("superset.tasks.scheduler.AsyncExecuteReportScheduleCommand") as
command,
+ patch("superset.tasks.scheduler.execute.update_state") as update_state,
+ ):
+ command.return_value.run.side_effect = error
+ execute(1234)
+
+ if attachment == "screenshot":
+ update_state.assert_not_called()
+ assert any(record.levelno == logging.WARNING for record in
caplog.records)
+ assert not any(record.levelno >= logging.ERROR for record in
caplog.records)
+ else:
+ update_state.assert_called_once_with(state="FAILURE")
+ assert any(
+ record.name == "superset.tasks.scheduler"
+ and record.levelno >= logging.ERROR
+ for record in caplog.records
+ )
diff --git a/tests/unit_tests/utils/csv_tests.py
b/tests/unit_tests/utils/csv_tests.py
index bac067874c7..3a3210a5171 100644
--- a/tests/unit_tests/utils/csv_tests.py
+++ b/tests/unit_tests/utils/csv_tests.py
@@ -21,8 +21,9 @@ from unittest import mock
import pandas as pd
import pyarrow as pa
-import pytest # noqa: F401
+import pytest
from pandas.api.types import is_datetime64_any_dtype
+from pytest_mock import MockerFixture
from superset.utils import csv, json
from superset.utils.core import GenericDataType
@@ -450,3 +451,27 @@ def test_get_chart_dataframe_forwards_timeout(monkeypatch:
pytest.MonkeyPatch) -
monkeypatch.setattr(csv, "get_chart_csv_data", fake)
get_chart_dataframe("http://dummy-url", timeout=99)
assert captured["timeout"] == 99
+
+
[email protected]("fail_read", [False, True])
+def test_get_chart_csv_data_closes_response(
+ mocker: MockerFixture, fail_read: bool
+) -> None:
+ """Release the legacy export connection on success and before retrying a
read."""
+ from superset.utils.csv import get_chart_csv_data
+
+ response = mocker.patch(
+ "urllib.request.build_opener"
+ ).return_value.open.return_value
+ response.getcode.return_value = 200
+ response.read.return_value = b"csv"
+ if fail_read:
+ response.read.side_effect = TimeoutError()
+ with pytest.raises(TimeoutError):
+ get_chart_csv_data("https://localhost/chart", {"session":
"cookie"}, 10)
+ else:
+ assert (
+ get_chart_csv_data("https://localhost/chart", {"session":
"cookie"}, 10)
+ == b"csv"
+ )
+ response.close.assert_called_once()