This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new b3674d43be0 `SparkSubmitHook`: Mask `_mask_cmd` secrets in linear time
(#71711)
b3674d43be0 is described below
commit b3674d43be01255430c0a673e7cc3b0a391e175e
Author: Divyanshu Singh <[email protected]>
AuthorDate: Sun Oct 4 22:13:11 2026 +0530
`SparkSubmitHook`: Mask `_mask_cmd` secrets in linear time (#71711)
* fix(spark): replace ReDoS-vulnerable regex in _mask_cmd with linear-time
token parser
The original regex in SparkSubmitHook._mask_cmd() used nested lazy
quantifiers with a lookahead that caused quadratic backtracking on
large inputs without secret/password keywords (O(n^2) — 57s for 50k
chars). Replace with a token-based approach that splits by whitespace
and masks values following sensitive keys in a single forward pass,
guaranteeing O(n) execution time (<1ms for 100k chars).
All existing test cases pass with identical output.
Closes #70676
* Add unit tests for SparkSubmitHook password masking
The masking rewrite changed how sensitive values are located in the command
string, so the space-separated, dotted-key, uppercase, multi-token quoted,
and
missing-value forms each need coverage. The large-input case guards against
a
regression to the quadratic behaviour that motivated the change.
* Mask spark-submit secrets separated by tabs or repeated spaces
Splitting the command on a single literal space missed values separated by
a tab
and consumed the empty field produced by a run of spaces as if it were the
value,
leaving the real secret in the output. Log lines fed through the same
masking are
column-aligned, so both shapes occur in practice.
Splitting on runs of whitespace and keeping the separators also preserves
the
original spacing rather than normalising it.
* Keep _mask_cmd an instance method
Converting it to a static method is unrelated to the masking change and only
widens the diff.
* Mask spark-submit secrets with a single linear-time pattern
Anchoring the match at a token boundary keeps the scan linear, so the
hand-rolled
tokeniser is no longer needed to avoid the backtracking. Matching the value
with
explicit quote alternatives, where a quote closes the value only when
whitespace
follows, keeps the previous masking behaviour intact.
* Describe _mask_cmd cost accurately rather than claiming linear time
Anchoring removes the retry-at-every-offset factor but a token packing many
sensitive keywords still backtracks quadratically, so the comments, test
name and
title should not promise O(n). Adds a regression test for that input shape.
* Stop masked quoted values from spanning newlines
The negated character classes for quoted values also matched newlines, so
an unterminated quote in a multiline argument or captured log line consumed
everything up to the next quote, swallowing unrelated output into the mask.
Exclude newlines from those classes and cover it with a regression test.
Also pass lists rather than tuples to _mask_cmd in the two performance
tests, which is what its signature declares and what mypy flagged in CI.
Co-Authored-By: Claude Opus 5 <[email protected]>
* refactor: tokenise _mask_cmd to fix multi-key blind spot
Switch from the anchored single-regex approach to tokenisation:
split on whitespace, run an unanchored regex per token, and handle
multi-token quoted values with manual lookahead.
This fixes the blind spot where a second sensitive key inside the
same whitespace-delimited token was missed by the (?<!\S) anchor
(e.g. Config(secret="x",password=hunter2) left hunter2 in the
clear).
Changes:
- _SENSITIVE_VALUE_RE -> _SENSITIVE_KV_RE (unanchored, simplified)
- _mask_sensitive_value -> _mask_sensitive_kv
- New _find_open_quote_key() helper for spanning quoted values
- _mask_cmd now tokenises, handles key=value / space-separated /
multi-token quoted values / unterminated quotes at newlines
- isinstance guard for plain-string connection_cmd (pre-existing bug)
- re.Match -> re.Match[str] type annotation
- Two new test cases pinning the reviewer-identified payloads
Drafted-by: Claude Code; reviewed by @divyanshus2404 before posting
* style: fix static check formatting errors
* Mask Spark submit secrets exactly as before while scanning in linear time
The token-based masking introduced earlier in this PR approximated the
original regex and diverged from it: a quoted value whose closing quote
was followed by punctuation, as in Python-repr or dict-shaped log output,
was left partly in the clear. Its per-token pattern also still backtracked
super-linearly on tokens with an "=" ahead of many keywords.
Scanning directly for what the original regex would match keeps the
masked output identical to main for every input while bounding the work
to a single pass, so long spark-submit log lines can no longer stall the
worker and no masking case is lost.
* fix(spark): extend quote limit to punctuation-terminated tokens
The previous _QUOTED_VALUE_LIMIT_RE only recognised a closing quote
when followed by whitespace ('(?=\s)'). Python-repr and dict-shaped
log lines that _submit_log_tail feeds through the masker routinely
produce tokens like Config(password="x", user=y) or {'password': 'a b'}
where the closing quote is immediately followed by a comma or bracket.
Change the lookahead from '(?=\s)' to '(?=\W|$)' so the limit fires on
any non-word character (whitespace, comma, paren, brace, ...) or end of
string. All pre-existing test_masks_passwords cases continue to pass;
two new cases that previously leaked (identified in potiuk's review)
now mask correctly.
* Restore whitespace-only quote limit when masking Spark secrets
Closing a quoted value at any quote followed by punctuation left the rest
of secrets such as 'Pa'$$w0rd' in the clear, diverging from the previous
masking behaviour, and made tokens of repeated closed quoted values scan
quadratically. The punctuation-terminated shapes it targeted were already
masked correctly by the whitespace-only limit.
---------
Co-authored-by: Divyanshu Singh <[email protected]>
Co-authored-by: divyanshus2404 <[email protected]>
Co-authored-by: Divyanshu Singh
<[email protected]>
Co-authored-by: Claude Opus 5 <[email protected]>
Co-authored-by: divyanshus2404 <[email protected]>
---
.../providers/apache/spark/hooks/spark_submit.py | 98 +++++++++++++----
.../unit/apache/spark/hooks/test_spark_submit.py | 118 +++++++++++++++++++++
2 files changed, 196 insertions(+), 20 deletions(-)
diff --git
a/providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py
b/providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py
index 87d838333ba..a0a65ccb11a 100644
---
a/providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py
+++
b/providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py
@@ -55,6 +55,81 @@ ALLOWED_SPARK_BINARIES = [DEFAULT_SPARK_BINARY,
"spark2-submit", "spark3-submit"
_K8S_WAIT_APP_COMPLETION_CONF = "spark.kubernetes.submission.waitAppCompletion"
+_SENSITIVE_KEYWORD_RE = re.compile(r"secret|password", re.IGNORECASE)
+_WHITESPACE_RE = re.compile(r"\s")
+_NON_WHITESPACE_RE = re.compile(r"\S")
+# Where a quoted value may stop at the latest: the quote followed by
whitespace, or a newline.
+_QUOTED_VALUE_LIMIT_RE = {quote: re.compile(rf"\n|{quote}(?=\s)") for quote in
("'", '"')}
+
+
+def _mask_sensitive_values(text: str) -> str:
+ r"""
+ Mask the value of every ``key=value`` / ``key value`` pair whose key
contains ``secret`` or ``password``.
+
+ Produces the same output as the single regular expression used previously::
+
+ (\S*?(?:secret|password)\S*?(?:=|\s+)(['"]?))(?:(?!\2\s).)*(\2) ->
\1******\3
+
+ but scans the input in linear time. That pattern backtracked quadratically
or worse on long tokens,
+ and it runs over arbitrary spark-submit output, so a single long log line
could stall the worker.
+
+ - The key starts where scanning resumed within the current token and ends
at the first ``=``
+ after the keyword, or at the whitespace ending the token.
+ - A value opening with a quote extends to the last matching quote before
either that quote
+ followed by whitespace or a newline; the value is then masked between
the quotes.
+ - Any other value, including an unterminated quoted one, is masked up to
the next whitespace.
+ """
+ length = len(text)
+ masked: list[str] = []
+ copied = 0
+ pos = 0
+ while True:
+ token = _NON_WHITESPACE_RE.search(text, pos)
+ if token is None:
+ break
+ token_start = token.start()
+ token_end_match = _WHITESPACE_RE.search(text, token_start)
+ token_end = token_end_match.start() if token_end_match else length
+ keyword = _SENSITIVE_KEYWORD_RE.search(text, token_start, token_end)
+ if keyword is None:
+ pos = token_end
+ continue
+ equals = text.find("=", keyword.end(), token_end)
+ if equals != -1:
+ value_start = equals + 1
+ elif token_end < length:
+ next_token = _NON_WHITESPACE_RE.search(text, token_end)
+ value_start = next_token.start() if next_token else length
+ else:
+ # Last token and nothing separates the key from a value.
+ break
+
+ if value_start < length and text[value_start] in "'\"":
+ quote = text[value_start]
+ limit = _QUOTED_VALUE_LIMIT_RE[quote].search(text, value_start + 1)
+ if limit is None:
+ limit_end = length
+ elif limit.group() == quote:
+ limit_end = limit.end()
+ else:
+ limit_end = limit.start()
+ closing = text.rfind(quote, value_start + 1, limit_end)
+ if closing != -1:
+ masked.append(text[copied : value_start + 1])
+ masked.append("******")
+ copied = closing
+ pos = closing + 1
+ continue
+
+ value_end_match = _WHITESPACE_RE.search(text, value_start)
+ value_end = value_end_match.start() if value_end_match else length
+ masked.append(text[copied:value_start])
+ masked.append("******")
+ copied = pos = value_end
+ masked.append(text[copied:])
+ return "".join(masked)
+
+
# The JVM's default uncaught-exception handler always prints this exact shape.
_EXCEPTION_START_RE = re.compile(r'Exception in thread "[^"]*"')
@@ -516,26 +591,9 @@ class SparkSubmitHook(BaseHook, LoggingMixin):
def _mask_cmd(self, connection_cmd: str | list[str]) -> str:
# Mask any password related fields in application args with key value
pair
# where key contains password (case insensitive), e.g.
HivePassword='abc'
- connection_cmd_masked = re.sub(
- r"("
- r"\S*?" # Match all non-whitespace characters before...
- r"(?:secret|password)" # ...literally a "secret" or "password"
- # word (not capturing them).
- r"\S*?" # All non-whitespace characters before either...
- r"(?:=|\s+)" # ...an equal sign or whitespace characters
- # (not capturing them).
- r"(['\"]?)" # An optional single or double quote.
- r")" # This is the end of the first capturing group.
- r"(?:(?!\2\s).)*" # All characters between optional quotes
- # (matched above); if the value is quoted,
- # it may contain whitespace.
- r"(\2)", # Optional matching quote.
- r"\1******\3",
- " ".join(connection_cmd),
- flags=re.I,
- )
-
- return connection_cmd_masked
+ if isinstance(connection_cmd, str):
+ connection_cmd = [connection_cmd]
+ return _mask_sensitive_values(" ".join(connection_cmd))
@property
def _submit_log_tail(self) -> str:
diff --git
a/providers/apache/spark/tests/unit/apache/spark/hooks/test_spark_submit.py
b/providers/apache/spark/tests/unit/apache/spark/hooks/test_spark_submit.py
index 0347c473160..52060da0b35 100644
--- a/providers/apache/spark/tests/unit/apache/spark/hooks/test_spark_submit.py
+++ b/providers/apache/spark/tests/unit/apache/spark/hooks/test_spark_submit.py
@@ -19,6 +19,7 @@ from __future__ import annotations
import base64
import os
+import time
from io import StringIO
from pathlib import Path
from types import ModuleType
@@ -1302,6 +1303,75 @@ class TestSparkSubmitHook:
("spark-submit",),
"spark-submit",
),
+ (
+ ("spark-submit", "foo", "--secret", "topsecret", "--bar"),
+ "spark-submit foo --secret ****** --bar",
+ ),
+ (
+ ("spark-submit", "--conf", "spark.mySecret=abc"),
+ "spark-submit --conf spark.mySecret=******",
+ ),
+ (
+ ("spark-submit", "--PASSWORD=abc"),
+ "spark-submit --PASSWORD=******",
+ ),
+ (
+ ("spark-submit", "--conf", "HivePassword='multi word pass'",
"--after"),
+ "spark-submit --conf HivePassword='******' --after",
+ ),
+ (
+ ("spark-submit", "--password", "'multi word pass'", "--bar",
"baz"),
+ "spark-submit --password '******' --bar baz",
+ ),
+ (
+ ("spark-submit", "--password"),
+ "spark-submit --password",
+ ),
+ (
+ ("spark-submit", "--password", "", "hunter2"),
+ "spark-submit --password ******",
+ ),
+ (
+ ("Using password hunter2",),
+ "Using password ******",
+ ),
+ (
+ ("spark-submit --password\thunter2",),
+ "spark-submit --password\t******",
+ ),
+ (
+ ("spark-submit\t--conf\tHivePassword='abc'",),
+ "spark-submit\t--conf\tHivePassword='******'",
+ ),
+ # Multiple sensitive keys inside a single token
(reviewer-identified blind spot):
+ # the old anchored regex missed the second key after a closed
quote.
+ (
+ ['Config(secret="x",password=hunter2)'],
+ 'Config(secret="******",password=******',
+ ),
+ (
+ ["--conf", "spark.a.secret='x',spark.b.password=hunter2"],
+ "--conf spark.a.secret='******',spark.b.password=******",
+ ),
+ # Quoted multi-word values whose closing quote is followed by
punctuation,
+ # as in Python-repr or dict-shaped log output.
+ (
+ ['Config(password="my pass word", user=x)'],
+ 'Config(password="******", user=x)',
+ ),
+ (
+ ["{'password': 'a b'}"],
+ "{'password': '******'}",
+ ),
+ # A quote followed by punctuation inside the value must not close
it early.
+ (
+ ["spark-submit", "--password='Pa'$$w0rd'", "--next"],
+ "spark-submit --password='******' --next",
+ ),
+ (
+ "spark-submit --password=hunter2",
+ "spark-submit --password=******",
+ ),
],
)
@pytest.mark.db_test
@@ -1315,6 +1385,54 @@ class TestSparkSubmitHook:
# Then
assert command_masked == expected
+ @pytest.mark.db_test
+ def test_masks_passwords_stays_fast_on_large_input(self) -> None:
+ # The previous pattern retried at every offset on long inputs, taking
tens of
+ # seconds for this payload and blocking the worker slot. The trailing
space is
+ # deliberate: it makes 25,000 separate tokens, which is what exercises
the
+ # per-offset retry rather than a single very long token.
+ hook = SparkSubmitHook()
+ payload = ["spark-submit", "--arg", "x " * 25_000]
+
+ start = time.monotonic()
+ command_masked = hook._mask_cmd(payload)
+ elapsed = time.monotonic() - start
+
+ assert command_masked == " ".join(payload)
+ assert elapsed < 5
+
+ @pytest.mark.db_test
+ def test_masks_passwords_does_not_swallow_following_lines(self) -> None:
+ # An unterminated quote must not consume the log lines after it: the
value
+ # ends at the newline, so the rest of the captured output survives
masking.
+ hook = SparkSubmitHook()
+ command = 'spark-submit --conf password="abc\nERROR: job
failed\n--other=1 "tail'
+
+ command_masked = hook._mask_cmd([command])
+
+ assert command_masked == 'spark-submit --conf password=******\nERROR:
job failed\n--other=1 "tail'
+
+ @pytest.mark.db_test
+ @pytest.mark.parametrize(
+ "token",
+ [
+ pytest.param("secret" * 20_000, id="repeated-keywords"),
+ pytest.param("a=" + "secret" * 20_000,
id="equals-before-repeated-keywords"),
+ pytest.param("password='x'," * 20_000,
id="repeated-closed-quoted-values"),
+ ],
+ )
+ def test_masks_passwords_stays_fast_on_repeated_keywords(self, token: str)
-> None:
+ # A token packing many sensitive keywords made the previous pattern
backtrack
+ # quadratically or worse; the scan must stay linear on these shapes.
+ hook = SparkSubmitHook()
+ payload = ["spark-submit", "--arg", token]
+
+ start = time.monotonic()
+ hook._mask_cmd(payload)
+ elapsed = time.monotonic() - start
+
+ assert elapsed < 5
+
@pytest.mark.db_test
def test_submit_log_tail_empty_when_no_lines_captured(self) -> None:
hook = SparkSubmitHook()